diff --git a/.gitignore b/.gitignore
index 38f80bb..57c1735 100644
--- a/.gitignore
+++ b/.gitignore
@@ -6,8 +6,10 @@ config.json
# Lokales Cold Backup (mbox) nicht versionieren
backup/
+# Gebaute Binaries (auch datierte Build-Varianten wie mail-graveyard-fix-identity-20260716)
/mail-graveyard
/mail-graveyard.exe
+/mail-graveyard-*
# Bun / Frontend
frontend-js/node_modules/
diff --git a/backend/00-router.go b/backend/00-router.go
index 2ab5beb..821fc3a 100644
--- a/backend/00-router.go
+++ b/backend/00-router.go
@@ -1305,7 +1305,7 @@ func renderTransferPreview(kind, value, folder string, index int, uid uint32) st
decodeHeader(msg.Header.Get("Subject")),
decodeHeader(msg.Header.Get("From")),
decodeHeader(msg.Header.Get("Date")),
- messageBody(raw),
+ messageBodyWithContext(raw, fmt.Sprintf("transfer kind=%q value=%q folder=%q index=%d uid=%d", kind, value, folder, index, uid)),
)
return b.String()
}
@@ -1621,7 +1621,7 @@ func renderSourceMailboxReadMessage(account Account, folder string, raw RawMessa
decodeHeader(msg.Header.Get("Subject")),
decodeHeader(msg.Header.Get("From")),
decodeHeader(msg.Header.Get("Date")),
- messageBody(raw.Body),
+ messageBodyWithContext(raw.Body, fmt.Sprintf("source account=%q folder=%q", account.Name, folder)),
)
return b.String()
}
@@ -1809,7 +1809,7 @@ func renderTargetMailboxReadMessage(account Account, folder string, raw RawMessa
decodeHeader(msg.Header.Get("Subject")),
decodeHeader(msg.Header.Get("From")),
decodeHeader(msg.Header.Get("Date")),
- messageBody(raw.Body),
+ messageBodyWithContext(raw.Body, fmt.Sprintf("target account=%q folder=%q", account.Name, folder)),
)
return b.String()
}
diff --git a/backend/02-database.go b/backend/02-database.go
index fd0923d..b972cdd 100644
--- a/backend/02-database.go
+++ b/backend/02-database.go
@@ -75,7 +75,13 @@ func ConnectDB(initDB bool) error {
if err != nil {
return err
}
- if _, err := db.Exec(`PRAGMA journal_mode=WAL; PRAGMA foreign_keys=ON;`); err != nil {
+ // SQLite pragmas are connection-local. Keep a single process-local
+ // connection so every query uses the busy timeout configured here. The web
+ // app and watcher are separate processes; WAL + busy_timeout coordinates
+ // their writes without failing immediately with SQLITE_BUSY.
+ db.SetMaxOpenConns(1)
+ db.SetMaxIdleConns(1)
+ if _, err := db.Exec(`PRAGMA busy_timeout=10000; PRAGMA journal_mode=WAL; PRAGMA foreign_keys=ON;`); err != nil {
_ = db.Close()
return err
}
@@ -115,14 +121,18 @@ func ConnectDB(initDB bool) error {
account_id INTEGER NOT NULL REFERENCES accounts(id) ON DELETE CASCADE,
folder TEXT NOT NULL,
message_id TEXT NOT NULL,
+ body_sha256 TEXT NOT NULL DEFAULT '',
+ mbox_done INTEGER NOT NULL DEFAULT 0,
+ target_done INTEGER NOT NULL DEFAULT 0,
copied_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
- UNIQUE(account_id, folder, message_id)
+ UNIQUE(account_id, folder, message_id, body_sha256)
)`,
`CREATE TABLE IF NOT EXISTS mbox_index(
account_id INTEGER NOT NULL,
folder TEXT NOT NULL,
seq INTEGER NOT NULL,
message_id TEXT NOT NULL,
+ body_sha256 TEXT NOT NULL DEFAULT '',
subject TEXT NOT NULL DEFAULT '',
from_addr TEXT NOT NULL DEFAULT '',
date TEXT NOT NULL DEFAULT '',
@@ -130,7 +140,7 @@ func ConnectDB(initDB bool) error {
frame_len INTEGER NOT NULL,
inner_offset INTEGER NOT NULL DEFAULT 0,
inner_len INTEGER NOT NULL,
- UNIQUE(account_id, folder, message_id)
+ UNIQUE(account_id, folder, seq)
)`,
`CREATE TABLE IF NOT EXISTS mbox_index_state(
account_id INTEGER NOT NULL,
@@ -183,6 +193,14 @@ func ConnectDB(initDB bool) error {
_ = db.Close()
return err
}
+ if err := ensureCopiedStageColumns(); err != nil {
+ _ = db.Close()
+ return err
+ }
+ if err := ensureMessageIdentitySchema(); err != nil {
+ _ = db.Close()
+ return err
+ }
if err := seedArchiveMailboxesFromFS(); err != nil {
_ = db.Close()
return err
@@ -190,6 +208,153 @@ func ConnectDB(initDB bool) error {
return nil
}
+// ensureMessageIdentitySchema upgrades the old Message-ID-only bookkeeping in
+// one transaction. No logical rows are discarded: old keys remain as rows with
+// an empty body hash and continue to act as compatibility aliases. The mbox
+// index is deliberately invalidated so its hashes are rebuilt from the archive
+// before the next migration pass makes a copy decision.
+func ensureMessageIdentitySchema() error {
+ copiedColumns, err := tableColumns("copied")
+ if err != nil {
+ return err
+ }
+ indexColumns, err := tableColumns("mbox_index")
+ if err != nil {
+ return err
+ }
+ upgradeCopied := !copiedColumns["body_sha256"]
+ upgradeIndex := !indexColumns["body_sha256"]
+ if !upgradeCopied && !upgradeIndex {
+ _, err := DB.Exec(`CREATE INDEX IF NOT EXISTS idx_copied_identity ON copied(account_id, folder, message_id, body_sha256);
+ CREATE INDEX IF NOT EXISTS idx_mbox_index_identity ON mbox_index(account_id, folder, message_id, body_sha256);`)
+ return err
+ }
+ tx, err := DB.Begin()
+ if err != nil {
+ return err
+ }
+ defer tx.Rollback()
+ if upgradeCopied {
+ statements := []string{
+ `ALTER TABLE copied RENAME TO copied_pre_identity`,
+ `CREATE TABLE copied(
+ account_id INTEGER NOT NULL REFERENCES accounts(id) ON DELETE CASCADE,
+ folder TEXT NOT NULL,
+ message_id TEXT NOT NULL,
+ body_sha256 TEXT NOT NULL DEFAULT '',
+ mbox_done INTEGER NOT NULL DEFAULT 0,
+ target_done INTEGER NOT NULL DEFAULT 0,
+ copied_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ UNIQUE(account_id, folder, message_id, body_sha256)
+ )`,
+ `INSERT INTO copied(account_id, folder, message_id, body_sha256, mbox_done, target_done, copied_at)
+ SELECT account_id, folder, message_id, '', mbox_done, target_done, copied_at FROM copied_pre_identity`,
+ `DROP TABLE copied_pre_identity`,
+ }
+ for _, stmt := range statements {
+ if _, err := tx.Exec(stmt); err != nil {
+ return err
+ }
+ }
+ }
+ if upgradeIndex {
+ statements := []string{
+ `ALTER TABLE mbox_index RENAME TO mbox_index_pre_identity`,
+ `CREATE TABLE mbox_index(
+ account_id INTEGER NOT NULL,
+ folder TEXT NOT NULL,
+ seq INTEGER NOT NULL,
+ message_id TEXT NOT NULL,
+ body_sha256 TEXT NOT NULL DEFAULT '',
+ subject TEXT NOT NULL DEFAULT '',
+ from_addr TEXT NOT NULL DEFAULT '',
+ date TEXT NOT NULL DEFAULT '',
+ file_offset INTEGER NOT NULL,
+ frame_len INTEGER NOT NULL,
+ inner_offset INTEGER NOT NULL DEFAULT 0,
+ inner_len INTEGER NOT NULL,
+ UNIQUE(account_id, folder, seq)
+ )`,
+ `INSERT INTO mbox_index(account_id, folder, seq, message_id, body_sha256, subject, from_addr, date, file_offset, frame_len, inner_offset, inner_len)
+ SELECT account_id, folder, seq, message_id, '', subject, from_addr, date, file_offset, frame_len, inner_offset, inner_len FROM mbox_index_pre_identity`,
+ `DROP TABLE mbox_index_pre_identity`,
+ `UPDATE mbox_index_state SET indexed_bytes=-1`,
+ }
+ for _, stmt := range statements {
+ if _, err := tx.Exec(stmt); err != nil {
+ return err
+ }
+ }
+ }
+ if _, err := tx.Exec(`CREATE INDEX IF NOT EXISTS idx_copied_identity ON copied(account_id, folder, message_id, body_sha256)`); err != nil {
+ return err
+ }
+ if _, err := tx.Exec(`CREATE INDEX IF NOT EXISTS idx_mbox_index_identity ON mbox_index(account_id, folder, message_id, body_sha256)`); err != nil {
+ return err
+ }
+ return tx.Commit()
+}
+
+func tableColumns(table string) (map[string]bool, error) {
+ rows, err := DB.Query(`PRAGMA table_info(` + table + `)`)
+ if err != nil {
+ return nil, err
+ }
+ defer rows.Close()
+ columns := map[string]bool{}
+ for rows.Next() {
+ var cid, notNull, pk int
+ var name, typ string
+ var defaultValue any
+ if err := rows.Scan(&cid, &name, &typ, ¬Null, &defaultValue, &pk); err != nil {
+ return nil, err
+ }
+ columns[name] = true
+ }
+ return columns, rows.Err()
+}
+
+func ensureCopiedStageColumns() error {
+ columns := map[string]bool{}
+ rows, err := DB.Query(`PRAGMA table_info(copied)`)
+ if err != nil {
+ return err
+ }
+ for rows.Next() {
+ var cid int
+ var name, typ string
+ var notNull int
+ var defaultValue any
+ var pk int
+ if err := rows.Scan(&cid, &name, &typ, ¬Null, &defaultValue, &pk); err != nil {
+ _ = rows.Close()
+ return err
+ }
+ columns[name] = true
+ }
+ if err := rows.Err(); err != nil {
+ _ = rows.Close()
+ return err
+ }
+ if err := rows.Close(); err != nil {
+ return err
+ }
+ // Existing copied rows were created by the old all-or-nothing pipeline and
+ // therefore represent both stages as complete. DEFAULT 1 preserves that
+ // truth while all new writes below set their stage values explicitly.
+ if !columns["mbox_done"] {
+ if _, err := DB.Exec(`ALTER TABLE copied ADD COLUMN mbox_done INTEGER NOT NULL DEFAULT 1`); err != nil {
+ return err
+ }
+ }
+ if !columns["target_done"] {
+ if _, err := DB.Exec(`ALTER TABLE copied ADD COLUMN target_done INTEGER NOT NULL DEFAULT 1`); err != nil {
+ return err
+ }
+ }
+ return nil
+}
+
func ensureDBFile(path string, initDB bool) error {
info, err := os.Stat(path)
if err == nil {
@@ -217,7 +382,6 @@ func ensureAppUserColumns() error {
if err != nil {
return err
}
- defer rows.Close()
for rows.Next() {
var cid int
var name, typ string
@@ -225,6 +389,7 @@ func ensureAppUserColumns() error {
var defaultValue any
var pk int
if err := rows.Scan(&cid, &name, &typ, ¬Null, &defaultValue, &pk); err != nil {
+ _ = rows.Close()
return err
}
if name == "display_name" {
@@ -232,6 +397,10 @@ func ensureAppUserColumns() error {
}
}
if err := rows.Err(); err != nil {
+ _ = rows.Close()
+ return err
+ }
+ if err := rows.Close(); err != nil {
return err
}
if !hasDisplayName {
@@ -518,22 +687,129 @@ func seedArchiveMailboxesFromFS() error {
}
func AlreadyCopied(accountID int64, folder, messageID string) (bool, error) {
- if messageID == "" {
- return false, nil
- }
- var x int
- err := DB.QueryRow(`SELECT 1 FROM copied WHERE account_id=? AND folder=? AND message_id=?`, accountID, folder, messageID).Scan(&x)
- if errors.Is(err, sql.ErrNoRows) {
- return false, nil
- }
- return err == nil, err
+ state, err := GetCopyState(accountID, folder, messageID)
+ return state.MboxDone && state.TargetDone, err
}
func MarkCopied(accountID int64, folder, messageID string) error {
- if messageID == "" {
+ return MarkIdentityCopied(accountID, folder, MessageIdentity{MessageID: messageID})
+}
+
+type CopyState struct {
+ MboxDone bool
+ TargetDone bool
+}
+
+func GetCopyState(accountID int64, folder, messageID string) (CopyState, error) {
+ return GetCopyIdentityState(accountID, folder, MessageIdentity{MessageID: messageID})
+}
+
+func MarkMboxCopied(accountID int64, folder, messageID string) error {
+ return MarkIdentityMboxCopied(accountID, folder, MessageIdentity{MessageID: messageID})
+}
+
+func MarkTargetCopied(accountID int64, folder, messageID string) error {
+ return MarkIdentityTargetCopied(accountID, folder, MessageIdentity{MessageID: messageID})
+}
+
+func GetCopyIdentityState(accountID int64, folder string, identity MessageIdentity) (CopyState, error) {
+ return GetCopyIdentityStateWithFolderAlias(accountID, folder, folder, identity)
+}
+
+func GetCopyIdentityStateWithFolderAlias(accountID int64, folder, legacyFolder string, identity MessageIdentity) (CopyState, error) {
+ if identity.MessageID == "" && identity.BodySHA256 == "" {
+ return CopyState{}, nil
+ }
+ folders := []string{folder}
+ if legacyFolder != "" && legacyFolder != folder {
+ folders = append(folders, legacyFolder)
+ }
+ // Precise rows always win. For Message-ID-less mail the hash is the whole
+ // identity, irrespective of which legacy key happened to be stored beside it.
+ var mboxDone, targetDone int
+ for _, candidateFolder := range folders {
+ var mbox, target int
+ var err error
+ if identity.MessageID == "" {
+ err = DB.QueryRow(`SELECT COALESCE(MAX(mbox_done),0), COALESCE(MAX(target_done),0) FROM copied
+ WHERE account_id=? AND folder=? AND body_sha256=?`, accountID, candidateFolder, identity.BodySHA256).Scan(&mbox, &target)
+ } else {
+ err = DB.QueryRow(`SELECT COALESCE(MAX(mbox_done),0), COALESCE(MAX(target_done),0) FROM copied
+ WHERE account_id=? AND folder=? AND message_id=? AND body_sha256=?`,
+ accountID, candidateFolder, identity.MessageID, identity.BodySHA256).Scan(&mbox, &target)
+ }
+ if err != nil {
+ return CopyState{}, err
+ }
+ mboxDone |= mbox
+ targetDone |= target
+ }
+ if mboxDone != 0 || targetDone != 0 {
+ return CopyState{MboxDone: mboxDone != 0, TargetDone: targetDone != 0}, nil
+ }
+ // Once a Message-ID has precise archive identities, a different body with
+ // that same ID is new and must not be hidden by the old wildcard row.
+ if identity.MessageID != "" && identity.BodySHA256 != "" {
+ var precise int
+ for _, candidateFolder := range folders {
+ var count int
+ if err := DB.QueryRow(`SELECT count(*) FROM copied WHERE account_id=? AND folder=? AND message_id=? AND body_sha256<>''`,
+ accountID, candidateFolder, identity.MessageID).Scan(&count); err != nil {
+ return CopyState{}, err
+ }
+ precise += count
+ }
+ if precise > 0 {
+ return CopyState{}, nil
+ }
+ }
+ aliases := []string{identity.MessageID, identity.LegacyMessageID}
+ if identity.MessageID == "" && identity.BodySHA256 != "" {
+ aliases = append(aliases, "sha256:"+identity.BodySHA256)
+ }
+ seen := map[string]bool{}
+ for _, alias := range aliases {
+ alias = strings.TrimSpace(alias)
+ if alias == "" || seen[alias] {
+ continue
+ }
+ seen[alias] = true
+ for _, candidateFolder := range folders {
+ var mbox, target int
+ err := DB.QueryRow(`SELECT COALESCE(MAX(mbox_done),0), COALESCE(MAX(target_done),0) FROM copied
+ WHERE account_id=? AND folder=? AND message_id=? AND body_sha256=''`, accountID, candidateFolder, alias).Scan(&mbox, &target)
+ if err != nil {
+ return CopyState{}, err
+ }
+ mboxDone |= mbox
+ targetDone |= target
+ }
+ }
+ return CopyState{MboxDone: mboxDone != 0, TargetDone: targetDone != 0}, nil
+}
+
+func MarkIdentityCopied(accountID int64, folder string, identity MessageIdentity) error {
+ return markCopyIdentity(accountID, folder, identity, true, true)
+}
+
+func MarkIdentityMboxCopied(accountID int64, folder string, identity MessageIdentity) error {
+ return markCopyIdentity(accountID, folder, identity, true, false)
+}
+
+func MarkIdentityTargetCopied(accountID int64, folder string, identity MessageIdentity) error {
+ return markCopyIdentity(accountID, folder, identity, false, true)
+}
+
+func markCopyIdentity(accountID int64, folder string, identity MessageIdentity, mboxDone, targetDone bool) error {
+ if identity.MessageID == "" && identity.BodySHA256 == "" {
return nil
}
- _, err := DB.Exec(`INSERT OR IGNORE INTO copied(account_id, folder, message_id) VALUES(?,?,?)`, accountID, folder, messageID)
+ _, err := DB.Exec(`INSERT INTO copied(account_id, folder, message_id, body_sha256, mbox_done, target_done)
+ VALUES(?,?,?,?,?,?)
+ ON CONFLICT(account_id, folder, message_id, body_sha256) DO UPDATE SET
+ mbox_done=MAX(copied.mbox_done, excluded.mbox_done),
+ target_done=MAX(copied.target_done, excluded.target_done),
+ copied_at=CURRENT_TIMESTAMP`, accountID, folder, identity.MessageID, identity.BodySHA256, boolInt(mboxDone), boolInt(targetDone))
return err
}
@@ -542,6 +818,7 @@ type MboxIndexEntry struct {
Folder string
Seq int
MessageID string
+ BodySHA256 string
Subject string
From string
Date string
@@ -552,7 +829,7 @@ type MboxIndexEntry struct {
}
func SaveMboxIndex(e MboxIndexEntry) error {
- if DB == nil || e.MessageID == "" {
+ if DB == nil || (e.MessageID == "" && e.BodySHA256 == "") {
return nil
}
var seq int
@@ -561,9 +838,9 @@ func SaveMboxIndex(e MboxIndexEntry) error {
} else {
_ = DB.QueryRow(`SELECT COALESCE(MAX(seq), -1) + 1 FROM mbox_index WHERE account_id=? AND folder=?`, e.AccountID, e.Folder).Scan(&seq)
}
- _, err := DB.Exec(`INSERT OR IGNORE INTO mbox_index(account_id, folder, seq, message_id, subject, from_addr, date, file_offset, frame_len, inner_offset, inner_len)
- VALUES(?,?,?,?,?,?,?,?,?,?,?)`,
- e.AccountID, e.Folder, seq, e.MessageID, e.Subject, e.From, e.Date, e.FileOffset, e.FrameLen, e.InnerOffset, e.InnerLen)
+ _, err := DB.Exec(`INSERT OR REPLACE INTO mbox_index(account_id, folder, seq, message_id, body_sha256, subject, from_addr, date, file_offset, frame_len, inner_offset, inner_len)
+ VALUES(?,?,?,?,?,?,?,?,?,?,?,?)`,
+ e.AccountID, e.Folder, seq, e.MessageID, e.BodySHA256, e.Subject, e.From, e.Date, e.FileOffset, e.FrameLen, e.InnerOffset, e.InnerLen)
return err
}
@@ -576,6 +853,53 @@ func ReplaceMboxIndex(accountID int64, folder string, entries []MboxIndexEntry,
return err
}
defer tx.Rollback()
+ folderAliases := []string{folder}
+ aliasRows, err := tx.Query(`SELECT folder FROM copied WHERE account_id=? UNION SELECT folder FROM mbox_index WHERE account_id=?`, accountID, accountID)
+ if err != nil {
+ return err
+ }
+ for aliasRows.Next() {
+ var candidate string
+ if err := aliasRows.Scan(&candidate); err != nil {
+ _ = aliasRows.Close()
+ return err
+ }
+ if candidate != folder && safeMboxName(candidate) == folder {
+ folderAliases = append(folderAliases, candidate)
+ }
+ }
+ if err := aliasRows.Err(); err != nil {
+ _ = aliasRows.Close()
+ return err
+ }
+ if err := aliasRows.Close(); err != nil {
+ return err
+ }
+ legacyIDsBySeq := map[int]string{}
+ for _, candidateFolder := range folderAliases {
+ rows, err := tx.Query(`SELECT seq, message_id FROM mbox_index WHERE account_id=? AND folder=?`, accountID, candidateFolder)
+ if err != nil {
+ return err
+ }
+ for rows.Next() {
+ var seq int
+ var messageID string
+ if err := rows.Scan(&seq, &messageID); err != nil {
+ _ = rows.Close()
+ return err
+ }
+ if _, exists := legacyIDsBySeq[seq]; !exists || candidateFolder == folder {
+ legacyIDsBySeq[seq] = messageID
+ }
+ }
+ if err := rows.Err(); err != nil {
+ _ = rows.Close()
+ return err
+ }
+ if err := rows.Close(); err != nil {
+ return err
+ }
+ }
if _, err := tx.Exec(`DELETE FROM mbox_index WHERE account_id=? AND folder=?`, accountID, folder); err != nil {
return err
}
@@ -583,15 +907,33 @@ func ReplaceMboxIndex(accountID int64, folder string, entries []MboxIndexEntry,
entry.AccountID = accountID
entry.Folder = folder
entry.Seq = i
- if entry.MessageID == "" {
+ if entry.MessageID == "" && entry.BodySHA256 == "" {
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,
+ if _, err := tx.Exec(`INSERT INTO mbox_index(account_id, folder, seq, message_id, body_sha256, subject, from_addr, date, file_offset, frame_len, inner_offset, inner_len)
+ VALUES(?,?,?,?,?,?,?,?,?,?,?,?)`,
+ entry.AccountID, entry.Folder, entry.Seq, entry.MessageID, entry.BodySHA256, entry.Subject, entry.From, entry.Date,
entry.FileOffset, entry.FrameLen, entry.InnerOffset, entry.InnerLen); err != nil {
return err
}
+ // A full rebuild recovers records which reached the old append-first
+ // pipeline but whose index/copied writes lost a SQLITE_BUSY race. In that
+ // pipeline the target append happened before the mbox append, so recovered
+ // records are complete in both stages. Marking them prevents a later watch
+ // run from appending duplicates to either destination.
+ state, found, err := copyStateForReindex(tx, entry.AccountID, folderAliases, entry, legacyIDsBySeq[i])
+ if err != nil {
+ return err
+ }
+ if !found {
+ state = CopyState{MboxDone: true, TargetDone: true}
+ }
+ if _, err := tx.Exec(`INSERT INTO copied(account_id, folder, message_id, body_sha256, mbox_done, target_done)
+ VALUES(?,?,?,?,?,?)
+ ON CONFLICT(account_id, folder, message_id, body_sha256) DO UPDATE SET mbox_done=1`,
+ entry.AccountID, entry.Folder, entry.MessageID, entry.BodySHA256, 1, boolInt(state.TargetDone)); err != nil {
+ return err
+ }
}
if _, err := tx.Exec(`INSERT INTO mbox_index_state(account_id, folder, indexed_bytes)
VALUES(?,?,?)
@@ -602,6 +944,88 @@ func ReplaceMboxIndex(accountID int64, folder string, entries []MboxIndexEntry,
return tx.Commit()
}
+// RemoveMboxIndexFolderAliases removes only superseded index metadata such as
+// "INBOX.Newbies " after the actual archive file has been rebuilt under its
+// canonical filesystem name "INBOX.Newbies". copied compatibility aliases and
+// all mbox files remain untouched.
+func RemoveMboxIndexFolderAliases(accountID int64, canonicalFolders map[string]bool) error {
+ rows, err := DB.Query(`SELECT folder FROM mbox_index WHERE account_id=? UNION SELECT folder FROM mbox_index_state WHERE account_id=?`, accountID, accountID)
+ if err != nil {
+ return err
+ }
+ var stale []string
+ for rows.Next() {
+ var folder string
+ if err := rows.Scan(&folder); err != nil {
+ _ = rows.Close()
+ return err
+ }
+ if !canonicalFolders[folder] && canonicalFolders[safeMboxName(folder)] {
+ stale = append(stale, folder)
+ }
+ }
+ if err := rows.Err(); err != nil {
+ _ = rows.Close()
+ return err
+ }
+ if err := rows.Close(); err != nil {
+ return err
+ }
+ if len(stale) == 0 {
+ return nil
+ }
+ tx, err := DB.Begin()
+ if err != nil {
+ return err
+ }
+ defer tx.Rollback()
+ for _, folder := range stale {
+ if _, err := tx.Exec(`DELETE FROM mbox_index WHERE account_id=? AND folder=?`, accountID, folder); err != nil {
+ return err
+ }
+ if _, err := tx.Exec(`DELETE FROM mbox_index_state WHERE account_id=? AND folder=?`, accountID, folder); err != nil {
+ return err
+ }
+ }
+ return tx.Commit()
+}
+
+func copyStateForReindex(tx *sql.Tx, accountID int64, folders []string, entry MboxIndexEntry, legacyMessageID string) (CopyState, bool, error) {
+ type candidate struct{ messageID, bodyHash string }
+ candidates := []candidate{{entry.MessageID, entry.BodySHA256}}
+ if legacyMessageID != "" {
+ candidates = append(candidates, candidate{legacyMessageID, ""})
+ }
+ if entry.MessageID != "" {
+ candidates = append(candidates, candidate{entry.MessageID, ""})
+ } else if entry.BodySHA256 != "" {
+ candidates = append(candidates, candidate{"sha256:" + entry.BodySHA256, ""})
+ }
+ seen := map[candidate]bool{}
+ var state CopyState
+ found := false
+ for _, c := range candidates {
+ if seen[c] {
+ continue
+ }
+ seen[c] = true
+ for _, folder := range folders {
+ var count, mboxDone, targetDone int
+ if err := tx.QueryRow(`SELECT count(*), COALESCE(MAX(mbox_done),0), COALESCE(MAX(target_done),0)
+ FROM copied WHERE account_id=? AND folder=? AND message_id=? AND body_sha256=?`,
+ accountID, folder, c.messageID, c.bodyHash).Scan(&count, &mboxDone, &targetDone); err != nil {
+ return CopyState{}, false, err
+ }
+ if count > 0 {
+ found = true
+ state.MboxDone = state.MboxDone || mboxDone != 0
+ state.TargetDone = state.TargetDone || targetDone != 0
+ }
+ }
+ }
+ return state, found, nil
+}
+
func UpdateMboxIndexState(accountID int64, folder string, indexedBytes int64) error {
if DB == nil || accountID == 0 {
return nil
@@ -629,7 +1053,7 @@ func ListMboxIndex(accountID int64, folder string) ([]MboxIndexEntry, error) {
if DB == nil || accountID == 0 {
return nil, sql.ErrNoRows
}
- rows, err := DB.Query(`SELECT account_id, folder, seq, message_id, subject, from_addr, date, file_offset, frame_len, inner_offset, inner_len
+ rows, err := DB.Query(`SELECT account_id, folder, seq, message_id, body_sha256, subject, from_addr, date, file_offset, frame_len, inner_offset, inner_len
FROM mbox_index WHERE account_id=? AND folder=? ORDER BY seq`, accountID, folder)
if err != nil {
return nil, err
@@ -638,7 +1062,7 @@ func ListMboxIndex(accountID int64, folder string) ([]MboxIndexEntry, error) {
var out []MboxIndexEntry
for rows.Next() {
var e MboxIndexEntry
- if err := rows.Scan(&e.AccountID, &e.Folder, &e.Seq, &e.MessageID, &e.Subject, &e.From, &e.Date, &e.FileOffset, &e.FrameLen, &e.InnerOffset, &e.InnerLen); err != nil {
+ if err := rows.Scan(&e.AccountID, &e.Folder, &e.Seq, &e.MessageID, &e.BodySHA256, &e.Subject, &e.From, &e.Date, &e.FileOffset, &e.FrameLen, &e.InnerOffset, &e.InnerLen); err != nil {
return nil, err
}
out = append(out, e)
@@ -650,10 +1074,10 @@ func GetMboxIndex(accountID int64, folder string, seq int) (MboxIndexEntry, erro
if DB == nil || accountID == 0 {
return MboxIndexEntry{}, sql.ErrNoRows
}
- row := DB.QueryRow(`SELECT account_id, folder, seq, message_id, subject, from_addr, date, file_offset, frame_len, inner_offset, inner_len
+ row := DB.QueryRow(`SELECT account_id, folder, seq, message_id, body_sha256, subject, from_addr, date, file_offset, frame_len, inner_offset, inner_len
FROM mbox_index WHERE account_id=? AND folder=? AND seq=?`, accountID, folder, seq)
var e MboxIndexEntry
- err := row.Scan(&e.AccountID, &e.Folder, &e.Seq, &e.MessageID, &e.Subject, &e.From, &e.Date, &e.FileOffset, &e.FrameLen, &e.InnerOffset, &e.InnerLen)
+ err := row.Scan(&e.AccountID, &e.Folder, &e.Seq, &e.MessageID, &e.BodySHA256, &e.Subject, &e.From, &e.Date, &e.FileOffset, &e.FrameLen, &e.InnerOffset, &e.InnerLen)
return e, err
}
diff --git a/backend/04-imap-source.go b/backend/04-imap-source.go
index 5b710a4..46daaad 100644
--- a/backend/04-imap-source.go
+++ b/backend/04-imap-source.go
@@ -1,6 +1,7 @@
package backend
import (
+ "bytes"
"crypto/sha256"
"crypto/tls"
"errors"
@@ -25,6 +26,17 @@ type RawMessage struct {
InternalDate time.Time // Original-Zeit, per APPEND erhalten
}
+// MessageIdentity is stable across the IMAP and mbox representations. A real
+// Message-ID is not sufficient on its own: broken mailers sometimes reuse it
+// for byte-different messages. Messages without a Message-ID are identified by
+// BodySHA256 alone. LegacyMessageID keeps the old raw-IMAP hash addressable
+// while existing databases are migrated without re-copying mail.
+type MessageIdentity struct {
+ MessageID string
+ BodySHA256 string
+ LegacyMessageID string
+}
+
type MessageHeader struct {
UID uint32
MessageID string
@@ -475,6 +487,43 @@ func messageID(body []byte) string {
return fmt.Sprintf("sha256:%x", sha256.Sum256(body))
}
+func identityForRawMessage(m RawMessage) MessageIdentity {
+ id := normalizeMessageID(m.MessageID)
+ legacyID := id
+ if strings.HasPrefix(strings.ToLower(id), "sha256:") {
+ id = ""
+ }
+ if id == "" {
+ if msg, err := mail.ReadMessage(bytes.NewReader(m.Body)); err == nil {
+ id = normalizeMessageID(msg.Header.Get("Message-ID"))
+ }
+ }
+ return MessageIdentity{
+ MessageID: id,
+ BodySHA256: bodySHA256(m.Body),
+ LegacyMessageID: legacyID,
+ }
+}
+
+// canonicalMessageBytes mirrors the bytes which can be recovered from the
+// current mbox writer/reader pair: line endings are LF, trailing record
+// separators are removed and mbox's >From escaping is undone. Hashing this
+// representation on both paths prevents the raw-IMAP/mbox hash split.
+func canonicalMessageBytes(raw []byte) []byte {
+ raw = bytes.ReplaceAll(raw, []byte("\r\n"), []byte("\n"))
+ lines := bytes.Split(raw, []byte("\n"))
+ for i, line := range lines {
+ if bytes.HasPrefix(line, []byte(">From ")) {
+ lines[i] = line[1:]
+ }
+ }
+ return bytes.TrimRight(bytes.Join(lines, []byte("\n")), "\n")
+}
+
+func bodySHA256(raw []byte) string {
+ return fmt.Sprintf("%x", sha256.Sum256(canonicalMessageBytes(raw)))
+}
+
func normalizeMessageID(id string) string {
return strings.Trim(strings.TrimSpace(id), "<>")
}
diff --git a/backend/04-imap-source_test.go b/backend/04-imap-source_test.go
index 1d05415..ef5d5ad 100644
--- a/backend/04-imap-source_test.go
+++ b/backend/04-imap-source_test.go
@@ -24,3 +24,34 @@ func TestMessageIDFallsBackToBodyHash(t *testing.T) {
t.Fatal("messageID hash fallback is not stable")
}
}
+
+func TestBodySHA256IsStableAcrossMboxNormalization(t *testing.T) {
+ raw := []byte("From: a@example.com\r\nSubject: Test\r\n\r\nFirst\r\nFrom escaped\r\n\r\n")
+ record := mboxRecord(RawMessage{Body: raw})
+ parts := splitMboxRecordsWithOffsets(record)
+ if len(parts) != 1 {
+ t.Fatalf("mbox parts=%d, want 1", len(parts))
+ }
+ if got, want := bodySHA256(parts[0].Message), bodySHA256(raw); got != want {
+ t.Fatalf("mbox hash=%s, raw hash=%s", got, want)
+ }
+}
+
+func TestIdentityUsesMessageIDAndBodyHash(t *testing.T) {
+ first := identityForRawMessage(RawMessage{MessageID: " Turkish İ test first
x",
+ want: []string{"Turkish İ test", "x"},
+ },
+ {
+ name: "kelvin sign",
+ in: "
x",
+ want: []string{"Kelvin K test", "x"},
+ },
+ {
+ name: "mixed blocks and repeated replacements",
+ in: "İ
third",
+ want: []string{"İ", "first", "K", "second", "third"},
+ drop: []string{"bad"},
+ },
+ }
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ got := htmlToText(tt.in)
+ for _, want := range tt.want {
+ if !strings.Contains(got, want) {
+ t.Fatalf("htmlToText(%q) = %q, missing %q", tt.in, got, want)
+ }
+ }
+ for _, drop := range tt.drop {
+ if strings.Contains(got, drop) {
+ t.Fatalf("htmlToText(%q) = %q, unexpectedly contains %q", tt.in, got, drop)
+ }
+ }
+ })
+ }
+}
+
+func TestASCIIFoldIndexReturnsOriginalByteOffset(t *testing.T) {
+ for _, tt := range []struct {
+ s string
+ needle string
+ want int
+ }{
+ {"İİ
x", "
", len("İİ")},
+ {"KK", "", len("KK")},
+ {"prefix
suffix", "
", len("prefix")},
+ } {
+ if got := asciiFoldIndex(tt.s, tt.needle); got != tt.want {
+ t.Fatalf("asciiFoldIndex(%q, %q) = %d, want %d", tt.s, tt.needle, got, tt.want)
+ }
+ }
+}
+
+func TestMessageBodyPanicFallsBackToRawMessage(t *testing.T) {
+ raw := []byte("Subject: Evidence\r\n\r\nraw body")
+ got := renderMessageBodySafely(raw, "test account=archive folder=INBOX seq=7", func() string {
+ panic("synthetic parser failure")
+ })
+ if !strings.Contains(got, "Darstellung fehlgeschlagen") || !strings.Contains(got, string(raw)) {
+ t.Fatalf("panic fallback did not preserve raw message: %q", got)
+ }
+}
+
func testRawMessage(id, subject string) RawMessage {
return RawMessage{
MessageID: id,
diff --git a/backend/07-migrate.go b/backend/07-migrate.go
index d866da5..97ccc59 100644
--- a/backend/07-migrate.go
+++ b/backend/07-migrate.go
@@ -3,6 +3,7 @@ package backend
import (
"fmt"
"log"
+ "os"
"path/filepath"
"strings"
"time"
@@ -16,8 +17,10 @@ import (
// Posteingang) landen im vorhandenen Rollen-Ordner des Ziels statt als
// Dublette; Pfad-Trenner wird umgehaengt. Dann EnsureFolder + Mails streamen.
// 3. Pro Mail: Message-ID gegen copied-Cache pruefen -> schon da? ueberspringen.
-// Sonst DOPPELT schreiben: TargetMailbox.Append (IMAP) + MboxWriter.Append
-// (lokal), dann MarkCopied.
+// Sonst zuerst ins lokale mbox-Sicherheitsarchiv schreiben und diesen
+// Zustand festhalten; erst danach folgt TargetMailbox.Append. Beide Stufen
+// sind getrennt wiederaufnehmbar, damit ein Zielausfall nie das Archiv
+// verhindert und ein Nachlauf keine mbox-Dublette erzeugt.
// 4. Fortschritt in jobs-Tabelle schreiben (Browser-Anzeige).
//
// name = Konto-Name oder "all". watch = true -> Delta-Schleife bis Abbruch
@@ -89,11 +92,6 @@ func migrateAccount(a Account, selectedFolders map[string]bool) error {
return fmt.Errorf("open source %s: %w", a.Name, err)
}
defer src.Close()
- dst, err := OpenIMAPTarget(a)
- if err != nil {
- return fmt.Errorf("open target %s: %w", a.Name, err)
- }
- defer dst.Close()
archiveName := accountMboxDir(a)
if archiveName == "" {
return fmt.Errorf("keine Archiv-Mailbox fuer %s ausgewaehlt", a.Name)
@@ -109,86 +107,187 @@ func migrateAccount(a Account, selectedFolders map[string]bool) error {
_ = finishJob(jobID, 0, 0, 1, "error")
return err
}
- targets, err := dst.Folders()
- if err != nil {
- _ = finishJob(jobID, 0, 0, 1, "error")
- return err
- }
folderMap, err := FolderMap(a.ID)
if err != nil {
_ = finishJob(jobID, 0, 0, 1, "error")
return err
}
+ var dst TargetMailbox
+ var targets []TargetFolder
+ if opened, openErr := OpenIMAPTarget(a); openErr != nil {
+ log.Printf("%s target unavailable; continuing archive-only: %v", a.Name, openErr)
+ } else {
+ dst = opened
+ targets, err = dst.Folders()
+ if err != nil {
+ log.Printf("%s target folders unavailable; continuing archive-only: %v", a.Name, err)
+ _ = dst.Close()
+ dst = nil
+ targets = nil
+ }
+ }
+ defer func() {
+ if dst != nil {
+ _ = dst.Close()
+ }
+ }()
total, done, errs := 0, 0, 0
+ archived, targeted, targetPending := 0, 0, 0
+ targetUnavailableLogged := false
for _, folder := range folders {
if len(selectedFolders) > 0 && !selectedFolders[folder.Name] {
continue
}
- dstFolder := MapSourceToTarget(folder.Name, folder.Attrs, folder.Delim, dst.Delim(), targets, folderMap[folder.Name])
+ dstDelim := "/"
+ if dst != nil {
+ dstDelim = dst.Delim()
+ }
+ dstFolder := MapSourceToTarget(folder.Name, folder.Attrs, folder.Delim, dstDelim, targets, folderMap[folder.Name])
if dstFolder == "" {
dstFolder = folder.Name
}
- if err := dst.EnsureFolder(dstFolder); err != nil {
- errs++
- log.Printf("%s %s ensure target %s error: %v", a.Name, folder.Name, dstFolder, err)
- continue
+ targetReady := false
+ if dst != nil {
+ if err := dst.EnsureFolder(dstFolder); err != nil {
+ log.Printf("%s %s ensure target %s error; retrying connection: %v", a.Name, folder.Name, dstFolder, err)
+ _ = dst.Close()
+ dst = nil
+ if reopened, reopenErr := OpenIMAPTarget(a); reopenErr != nil {
+ log.Printf("%s target reconnect failed; continuing archive-only: %v", a.Name, reopenErr)
+ } else if ensureErr := reopened.EnsureFolder(dstFolder); ensureErr != nil {
+ log.Printf("%s %s ensure target %s after reconnect failed; continuing archive-only: %v", a.Name, folder.Name, dstFolder, ensureErr)
+ _ = reopened.Close()
+ } else {
+ dst = reopened
+ targetReady = true
+ }
+ } else {
+ targetReady = true
+ }
}
folderDone, folderErrs := 0, 0
+ folderArchived, folderTargeted, folderPending := 0, 0, 0
folderTotal := 0
+ identityFolder := safeMboxName(folder.Name)
+ // After an identity-schema upgrade the indexed byte count is invalidated.
+ // Rebuild from the actual archive before consulting copied, otherwise a
+ // legacy Message-ID row could hide a byte-different message.
+ archivePath := filepath.Join(mboxDir, safeMboxName(folder.Name)+mboxFileExt())
+ if _, statErr := os.Stat(archivePath); statErr == nil {
+ if indexErr := EnsureMboxIndexed(archivePath); indexErr != nil {
+ errs++
+ folderErrs++
+ log.Printf("%s %s archive identity index error: %v", a.Name, folder.Name, indexErr)
+ continue
+ }
+ } else if !os.IsNotExist(statErr) {
+ errs++
+ folderErrs++
+ log.Printf("%s %s archive stat error: %v", a.Name, folder.Name, statErr)
+ continue
+ }
if err := src.Fetch(folder.Name, func(m RawMessage) error {
total++
folderTotal++
- already, err := AlreadyCopied(a.ID, folder.Name, m.MessageID)
+ identity := identityForRawMessage(m)
+ state, err := GetCopyIdentityStateWithFolderAlias(a.ID, identityFolder, folder.Name, identity)
if err != nil {
errs++
folderErrs++
log.Printf("%s %s dedup error %s: %v", a.Name, folder.Name, m.MessageID, err)
return nil
}
- if already {
+ if state.MboxDone && state.TargetDone {
return nil
}
- if err := dst.Append(dstFolder, m); err != nil {
- errs++
- folderErrs++
- log.Printf("%s %s append to %s error %s: %v", a.Name, folder.Name, dstFolder, m.MessageID, err)
- return nil
+ if !state.MboxDone {
+ appendInfo, err := mbox.Append(folder.Name, m)
+ if err != nil {
+ errs++
+ folderErrs++
+ log.Printf("%s %s mbox error %s: %v", a.Name, folder.Name, m.MessageID, err)
+ return nil
+ }
+ if err := SaveMboxIndex(MboxIndexEntry{
+ AccountID: a.ID,
+ Folder: identityFolder,
+ MessageID: identity.MessageID,
+ BodySHA256: identity.BodySHA256,
+ Subject: appendInfo.Subject,
+ From: appendInfo.From,
+ Date: appendInfo.Date,
+ FileOffset: appendInfo.FileOffset,
+ FrameLen: appendInfo.FrameLen,
+ InnerOffset: appendInfo.InnerOffset,
+ InnerLen: appendInfo.InnerLen,
+ }); err != nil {
+ errs++
+ folderErrs++
+ log.Printf("%s %s mbox index error %s: %v", a.Name, folder.Name, m.MessageID, err)
+ return nil
+ }
+ if err := UpdateMboxIndexState(a.ID, identityFolder, 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 := MarkIdentityMboxCopied(a.ID, identityFolder, identity); err != nil {
+ errs++
+ folderErrs++
+ log.Printf("%s %s mbox stage mark error %s: %v", a.Name, folder.Name, m.MessageID, err)
+ return nil
+ }
+ state.MboxDone = true
+ archived++
+ folderArchived++
}
- appendInfo, err := mbox.Append(folder.Name, m)
- if err != nil {
- errs++
- folderErrs++
- log.Printf("%s %s mbox error %s: %v", a.Name, folder.Name, m.MessageID, err)
- return nil
- }
- if err := SaveMboxIndex(MboxIndexEntry{
- AccountID: a.ID,
- Folder: folder.Name,
- MessageID: m.MessageID,
- Subject: appendInfo.Subject,
- From: appendInfo.From,
- Date: appendInfo.Date,
- FileOffset: appendInfo.FileOffset,
- FrameLen: appendInfo.FrameLen,
- InnerOffset: appendInfo.InnerOffset,
- InnerLen: appendInfo.InnerLen,
- }); err != nil {
- errs++
- folderErrs++
- 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++
- log.Printf("%s %s mark error %s: %v", a.Name, folder.Name, m.MessageID, err)
- return nil
+ if !state.TargetDone {
+ if dst == nil || !targetReady {
+ errs++
+ folderErrs++
+ targetPending++
+ folderPending++
+ if !targetUnavailableLogged {
+ log.Printf("%s target unavailable; messages remain safely archived with target_done=0", a.Name)
+ targetUnavailableLogged = true
+ }
+ return nil
+ }
+ if err := dst.Append(dstFolder, m); err != nil {
+ log.Printf("%s %s append to %s error %s; retrying connection: %v", a.Name, folder.Name, dstFolder, m.MessageID, err)
+ _ = dst.Close()
+ dst = nil
+ reopened, reopenErr := OpenIMAPTarget(a)
+ if reopenErr == nil {
+ reopenErr = reopened.EnsureFolder(dstFolder)
+ }
+ if reopenErr == nil {
+ reopenErr = reopened.Append(dstFolder, m)
+ }
+ if reopenErr != nil {
+ if reopened != nil {
+ _ = reopened.Close()
+ }
+ errs++
+ folderErrs++
+ targetPending++
+ folderPending++
+ targetReady = false
+ log.Printf("%s target reconnect/append failed; continuing archive-only: %v", a.Name, reopenErr)
+ return nil
+ }
+ dst = reopened
+ targetReady = true
+ }
+ if err := MarkIdentityTargetCopied(a.ID, identityFolder, identity); err != nil {
+ errs++
+ folderErrs++
+ log.Printf("%s %s target stage mark error %s: %v", a.Name, folder.Name, m.MessageID, err)
+ return nil
+ }
+ targeted++
+ folderTargeted++
}
done++
folderDone++
@@ -198,14 +297,16 @@ func migrateAccount(a Account, selectedFolders map[string]bool) error {
folderErrs++
log.Printf("%s %s fetch error: %v", a.Name, folder.Name, err)
}
- log.Printf("migration %s %s -> %s: total=%d copied=%d errors=%d", a.Name, folder.Name, dstFolder, folderTotal, folderDone, folderErrs)
+ log.Printf("migration %s %s -> %s: total=%d complete=%d archived=%d target=%d pending=%d errors=%d",
+ a.Name, folder.Name, dstFolder, folderTotal, folderDone, folderArchived, folderTargeted, folderPending, folderErrs)
}
state := "success"
if errs > 0 {
state = "partial"
}
_ = finishJob(jobID, total, done, errs, state)
- log.Printf("migration %s all folders: total=%d copied=%d errors=%d", a.Name, total, done, errs)
+ log.Printf("migration %s all folders: total=%d complete=%d archived=%d target=%d pending=%d errors=%d",
+ a.Name, total, done, archived, targeted, targetPending, errs)
if errs > 0 {
return fmt.Errorf("migration %s finished with %d errors", a.Name, errs)
}
diff --git a/backend/08-viewer.go b/backend/08-viewer.go
index e3fd2b4..1228693 100644
--- a/backend/08-viewer.go
+++ b/backend/08-viewer.go
@@ -62,7 +62,7 @@ func messageHandler(w http.ResponseWriter, r *http.Request) {
decodeHeader(msg.Header.Get("Subject")),
decodeHeader(msg.Header.Get("From")),
decodeHeader(msg.Header.Get("Date")),
- messageBody(raw),
+ messageBodyWithContext(raw, fmt.Sprintf("archive account=%q folder=%q seq=%d", account, folder, index)),
)
}
@@ -393,7 +393,7 @@ func renderTargetMatch(b *strings.Builder, match forwardedMatch) {
decodeHeader(msg.Header.Get("Subject")),
decodeHeader(msg.Header.Get("From")),
decodeHeader(msg.Header.Get("Date"))+" | "+match.Location(),
- messageBody(match.Raw),
+ messageBodyWithContext(match.Raw, "target "+match.Location()),
)
}
@@ -487,7 +487,8 @@ func filterEntries(path string, entries []MboxEntry, q string) []MboxEntry {
continue
}
raw, err := ReadMboxMessage(path, entry.Index)
- if err == nil && strings.Contains(strings.ToLower(messageBody(raw)), needle) {
+ if err == nil && strings.Contains(strings.ToLower(messageBodyWithContext(raw,
+ fmt.Sprintf("archive search path=%q seq=%d", path, entry.Index))), needle) {
out = append(out, entry)
}
}
@@ -523,7 +524,7 @@ func findTargetIMAPCopy(account, targetAccount string, raw []byte) (forwardedMat
if err != nil {
return forwardedMatch{}, false, err
}
- match, ok, scanErr := scanTargetForMessage(dst, targetUser, id)
+ match, ok, scanErr := scanTargetForMessage(dst, targetUser, id, bodySHA256(raw))
closeErr := dst.Close()
if scanErr != nil {
return forwardedMatch{}, false, scanErr
@@ -538,7 +539,7 @@ func findTargetIMAPCopy(account, targetAccount string, raw []byte) (forwardedMat
return forwardedMatch{}, false, nil
}
-func scanTargetForMessage(dst TargetMailbox, targetUser, messageID string) (forwardedMatch, bool, error) {
+func scanTargetForMessage(dst TargetMailbox, targetUser, messageID, bodyHash string) (forwardedMatch, bool, error) {
messageID = normalizeMessageID(messageID)
folders, err := dst.Folders()
if err != nil {
@@ -559,6 +560,9 @@ func scanTargetForMessage(dst TargetMailbox, targetUser, messageID string) (forw
if err != nil {
continue
}
+ if bodySHA256(msg.Body) != bodyHash {
+ continue
+ }
return forwardedMatch{
Account: "Ziel: " + targetUser,
Folder: folder.Name,
@@ -582,6 +586,7 @@ func findForwardedCopy(account, targetAccount string, raw []byte) (forwardedMatc
if err != nil {
return forwardedMatch{}, false
}
+ bodyHash := bodySHA256(raw)
for _, entry := range entries {
if !entry.IsDir() || entry.Name() == account {
continue
@@ -600,7 +605,7 @@ func findForwardedCopy(account, targetAccount string, raw []byte) (forwardedMatc
if err != nil {
continue
}
- if messageIDHeader(candidate) == id {
+ if messageIDHeader(candidate) == id && bodySHA256(candidate) == bodyHash {
return forwardedMatch{
Account: entry.Name(),
Folder: folderNameFromMboxPath(path),
@@ -713,7 +718,7 @@ func appendManualTargetCopy(account, folder, targetAccount string, raw []byte) (
return "", err
}
if a.ID != 0 {
- _ = MarkCopied(a.ID, folder, id)
+ _ = MarkIdentityCopied(a.ID, folder, identityForRawMessage(msg))
}
return "Manuell kopiert nach " + strings.TrimSpace(a.DstUser) + " / " + dstFolder + ".", nil
}
diff --git a/backend/13-target-dedup.go b/backend/13-target-dedup.go
index b1e3999..8790972 100644
--- a/backend/13-target-dedup.go
+++ b/backend/13-target-dedup.go
@@ -15,21 +15,24 @@ type TargetDedupReport struct {
ExtraCopies int
Candidates int
WithoutID int
+ VariantGroups int
Deleted int
FolderReports []TargetDedupFolderReport
}
type TargetDedupFolderReport struct {
- Folder string
- Groups int
- ExtraCopies int
- WithoutID int
- Deletes []TargetDedupDelete
+ Folder string
+ Groups int
+ ExtraCopies int
+ WithoutID int
+ VariantGroups int
+ Deletes []TargetDedupDelete
}
type TargetDedupDelete struct {
- UID uint32
- MessageID string
+ UID uint32
+ MessageID string
+ BodySHA256 string
}
func DedupTarget(name string, folders []string, apply bool) (TargetDedupReport, error) {
@@ -61,13 +64,14 @@ func DedupTarget(name string, folders []string, apply bool) (TargetDedupReport,
report.ExtraCopies += folderReport.ExtraCopies
report.Candidates += len(folderReport.Deletes)
report.WithoutID += folderReport.WithoutID
+ report.VariantGroups += folderReport.VariantGroups
if apply {
report.Deleted += len(folderReport.Deletes)
}
report.FolderReports = append(report.FolderReports, folderReport)
}
- log.Printf("target dedup %s apply=%v folders=%d groups=%d extra=%d candidates=%d deleted=%d without_id=%d",
- report.Account, report.Apply, report.Folders, report.Groups, report.ExtraCopies, report.Candidates, report.Deleted, report.WithoutID)
+ log.Printf("target dedup %s apply=%v folders=%d groups=%d extra=%d candidates=%d deleted=%d without_id=%d protected_variant_groups=%d",
+ report.Account, report.Apply, report.Folders, report.Groups, report.ExtraCopies, report.Candidates, report.Deleted, report.WithoutID, report.VariantGroups)
return report, nil
}
@@ -91,12 +95,35 @@ func dedupTargetFolder(dst TargetMailbox, folder string, apply bool) (TargetDedu
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)
+ byBody := map[string][]uint32{}
+ for _, uid := range uids {
+ msg, err := dst.FetchOne(folder, uid)
+ if err != nil {
+ return report, fmt.Errorf("%s fetch uid %d for safe dedup: %w", folder, uid, err)
+ }
+ hash := bodySHA256(msg.Body)
+ byBody[hash] = append(byBody[hash], uid)
+ }
+ if len(byBody) > 1 {
+ report.VariantGroups++
+ }
+ hashes := make([]string, 0, len(byBody))
+ for hash := range byBody {
+ hashes = append(hashes, hash)
+ }
+ sort.Strings(hashes)
+ for _, hash := range hashes {
+ bodyUIDs := byBody[hash]
+ if len(bodyUIDs) < 2 {
+ continue
+ }
+ sort.Slice(bodyUIDs, func(i, j int) bool { return bodyUIDs[i] < bodyUIDs[j] })
+ report.Groups++
+ report.ExtraCopies += len(bodyUIDs) - 1
+ for _, uid := range bodyUIDs[1:] {
+ report.Deletes = append(report.Deletes, TargetDedupDelete{UID: uid, MessageID: id, BodySHA256: hash})
+ deleteUIDs = append(deleteUIDs, uid)
+ }
}
}
sort.Slice(report.Deletes, func(i, j int) bool {
@@ -108,14 +135,14 @@ func dedupTargetFolder(dst TargetMailbox, folder string, apply bool) (TargetDedu
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)
+ log.Printf("target dedup delete folder=%s uid=%d message_id=%s body_sha256=%s", folder, del.UID, del.MessageID, del.BodySHA256)
}
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)
+ if report.Groups > 0 || report.WithoutID > 0 || report.VariantGroups > 0 {
+ log.Printf("target dedup folder=%s groups=%d extra=%d without_id=%d protected_variant_groups=%d apply=%v", folder, report.Groups, report.ExtraCopies, report.WithoutID, report.VariantGroups, apply)
}
return report, nil
}
@@ -125,13 +152,13 @@ func LogTargetDedupReport(report TargetDedupReport) {
if report.Apply {
mode = "APPLY"
}
- log.Printf("target dedup report %s account=%s folders=%d groups=%d extra=%d candidates=%d deleted=%d without_id=%d",
- mode, report.Account, report.Folders, report.Groups, report.ExtraCopies, report.Candidates, report.Deleted, report.WithoutID)
+ log.Printf("target dedup report %s account=%s folders=%d groups=%d extra=%d candidates=%d deleted=%d without_id=%d protected_variant_groups=%d",
+ mode, report.Account, report.Folders, report.Groups, report.ExtraCopies, report.Candidates, report.Deleted, report.WithoutID, report.VariantGroups)
for _, folder := range report.FolderReports {
- if folder.Groups == 0 && folder.WithoutID == 0 {
+ if folder.Groups == 0 && folder.WithoutID == 0 && folder.VariantGroups == 0 {
continue
}
- log.Printf("target dedup report folder=%s groups=%d extra=%d without_id=%d", folder.Folder, folder.Groups, folder.ExtraCopies, folder.WithoutID)
+ log.Printf("target dedup report folder=%s groups=%d extra=%d without_id=%d protected_variant_groups=%d", folder.Folder, folder.Groups, folder.ExtraCopies, folder.WithoutID, folder.VariantGroups)
for _, del := range folder.Deletes {
log.Printf("target dedup report candidate folder=%s uid=%d message_id=%s", folder.Folder, del.UID, del.MessageID)
}
diff --git a/backend/13-target-dedup_test.go b/backend/13-target-dedup_test.go
new file mode 100644
index 0000000..d7ed764
--- /dev/null
+++ b/backend/13-target-dedup_test.go
@@ -0,0 +1,70 @@
+package backend
+
+import (
+ "fmt"
+ "reflect"
+ "testing"
+)
+
+type dedupTargetStub struct {
+ headers []MessageHeader
+ bodies map[uint32][]byte
+ deleted []uint32
+}
+
+func (s *dedupTargetStub) Folders() ([]TargetFolder, error) { return nil, nil }
+func (s *dedupTargetStub) Delim() string { return "/" }
+func (s *dedupTargetStub) EnsureFolder(string) error { return nil }
+func (s *dedupTargetStub) Append(string, RawMessage) error { return nil }
+func (s *dedupTargetStub) Headers(string, int, int) ([]MessageHeader, error) {
+ return s.headers, nil
+}
+func (s *dedupTargetStub) AllHeaders(string) ([]MessageHeader, error) { return s.headers, nil }
+func (s *dedupTargetStub) FetchOne(_ string, uid uint32) (RawMessage, error) {
+ body, ok := s.bodies[uid]
+ if !ok {
+ return RawMessage{}, fmt.Errorf("missing uid %d", uid)
+ }
+ return RawMessage{Body: body}, nil
+}
+func (s *dedupTargetStub) DeleteUIDs(_ string, uids []uint32) error {
+ s.deleted = append(s.deleted, uids...)
+ return nil
+}
+func (s *dedupTargetStub) Close() error { return nil }
+
+func TestDedupTargetFolderProtectsDifferentBodiesWithSameMessageID(t *testing.T) {
+ sameA := []byte("Message-ID:
`, `
Turkish İ test