diff --git a/backend/02-database.go b/backend/02-database.go index 4b976d7..4921033 100644 --- a/backend/02-database.go +++ b/backend/02-database.go @@ -2,6 +2,9 @@ package backend import ( "database/sql" + "encoding/json" + "errors" + "os" _ "modernc.org/sqlite" ) @@ -11,27 +14,27 @@ var DB *sql.DB // Account = eine Umzugs-Zeile: Quelle (alt) -> Ziel (neu), frei gemapptes // Zieladress-Postfach. Wird im Browser gepflegt (Tabelle accounts). type Account struct { - ID int64 - Name string // z.B. "hans-peter" + ID int64 `json:"id"` + Name string `json:"name"` // z.B. "hans-peter" // Quelle (alter Provider). Security: "tls" (implizit, 993/995), // "starttls" (Klartext-Port 143/110 + Upgrade), "none" (reiner Klartext). // SrcInsecure = ungueltige/selbstsignierte Zerts akzeptieren (alte Hoster). - SrcHost string - SrcPort int - SrcSecurity string - SrcInsecure bool - SrcUser string // hans-peter@dr-gold.de - SrcPass string - SrcProto string // "imap" (default) oder "pop3" (Fallback-Quelle) + SrcHost string `json:"src_host"` + SrcPort int `json:"src_port"` + SrcSecurity string `json:"src_security"` + SrcInsecure bool `json:"src_insecure"` + SrcUser string `json:"src_user"` // hans-peter@dr-gold.de + SrcPass string `json:"src_pass"` + SrcProto string `json:"src_proto"` // "imap" (default) oder "pop3" (Fallback-Quelle) // Ziel (neuer Provider). Meist "tls"; "none"/"starttls" ebenso moeglich. - DstHost string - DstPort int - DstSecurity string - DstInsecure bool - DstUser string // archiv-hans-peter@dr-gold.com - DstPass string - MboxDir string // Unterordner unter Cfg.MboxRoot; leer = Cfg.MboxRoot/Name - Active bool + DstHost string `json:"dst_host"` + DstPort int `json:"dst_port"` + DstSecurity string `json:"dst_security"` + DstInsecure bool `json:"dst_insecure"` + DstUser string `json:"dst_user"` // archiv-hans-peter@dr-gold.com + DstPass string `json:"dst_pass"` + MboxDir string `json:"mbox_dir"` // Unterordner unter Cfg.MboxRoot; leer = Cfg.MboxRoot/Name + Active bool `json:"active"` } // ConnectDB oeffnet die SQLite-DB (modernc, kein cgo) und legt das Schema an. @@ -49,11 +52,202 @@ type Account struct { // jobs(account_id, started, finished, -- Lauf-Historie fuer den Fortschritt // total, done, errors, state) im Browser. func ConnectDB() error { - // TODO Codex: sql.Open("sqlite", Cfg.DBPath), PRAGMA journal_mode=WAL, - // CREATE TABLE IF NOT EXISTS ... siehe Schema oben. + if Cfg.DBPath == "" { + Cfg.DBPath = "emailforwarder.db" + } + db, err := sql.Open("sqlite", Cfg.DBPath) + if err != nil { + return err + } + if _, err := db.Exec(`PRAGMA journal_mode=WAL; PRAGMA foreign_keys=ON;`); err != nil { + _ = db.Close() + return err + } + schema := []string{ + `CREATE TABLE IF NOT EXISTS accounts( + id INTEGER PRIMARY KEY AUTOINCREMENT, + name TEXT NOT NULL UNIQUE, + src_host TEXT NOT NULL, + src_port INTEGER NOT NULL, + src_security TEXT NOT NULL DEFAULT 'tls', + src_insecure INTEGER NOT NULL DEFAULT 0, + src_user TEXT NOT NULL, + src_pass TEXT NOT NULL, + src_proto TEXT NOT NULL DEFAULT 'imap', + dst_host TEXT NOT NULL, + dst_port INTEGER NOT NULL, + dst_security TEXT NOT NULL DEFAULT 'tls', + dst_insecure INTEGER NOT NULL DEFAULT 0, + dst_user TEXT NOT NULL, + dst_pass TEXT NOT NULL, + mbox_dir TEXT NOT NULL DEFAULT '', + active INTEGER NOT NULL DEFAULT 1 + )`, + `CREATE TABLE IF NOT EXISTS folder_map( + id INTEGER PRIMARY KEY AUTOINCREMENT, + account_id INTEGER NOT NULL REFERENCES accounts(id) ON DELETE CASCADE, + src_folder TEXT NOT NULL, + dst_folder TEXT NOT NULL, + UNIQUE(account_id, src_folder) + )`, + `CREATE TABLE IF NOT EXISTS copied( + account_id INTEGER NOT NULL REFERENCES accounts(id) ON DELETE CASCADE, + folder TEXT NOT NULL, + message_id TEXT NOT NULL, + copied_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP, + UNIQUE(account_id, folder, message_id) + )`, + `CREATE TABLE IF NOT EXISTS jobs( + id INTEGER PRIMARY KEY AUTOINCREMENT, + account_id INTEGER NOT NULL REFERENCES accounts(id) ON DELETE CASCADE, + started TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP, + finished TEXT, + total INTEGER NOT NULL DEFAULT 0, + done INTEGER NOT NULL DEFAULT 0, + errors INTEGER NOT NULL DEFAULT 0, + state TEXT NOT NULL DEFAULT 'running' + )`, + } + for _, stmt := range schema { + if _, err := db.Exec(stmt); err != nil { + _ = db.Close() + return err + } + } + DB = db return nil } -// TODO Codex: ListAccounts / SaveAccount / DeleteAccount / GetAccount, -// AlreadyCopied(accountID, folder, messageID) bool, -// MarkCopied(accountID, folder, messageID). +func ListAccounts() ([]Account, error) { + rows, err := DB.Query(`SELECT id,name,src_host,src_port,src_security,src_insecure,src_user,src_pass,src_proto, + dst_host,dst_port,dst_security,dst_insecure,dst_user,dst_pass,mbox_dir,active + FROM accounts ORDER BY name`) + if err != nil { + return nil, err + } + defer rows.Close() + var out []Account + for rows.Next() { + a, err := scanAccount(rows) + if err != nil { + return nil, err + } + out = append(out, a) + } + return out, rows.Err() +} + +func GetAccount(name string) (Account, error) { + row := DB.QueryRow(`SELECT id,name,src_host,src_port,src_security,src_insecure,src_user,src_pass,src_proto, + dst_host,dst_port,dst_security,dst_insecure,dst_user,dst_pass,mbox_dir,active + FROM accounts WHERE name=?`, name) + return scanAccount(row) +} + +func SaveAccount(a Account) error { + normalizeAccount(&a) + _, err := DB.Exec(`INSERT INTO accounts(name,src_host,src_port,src_security,src_insecure,src_user,src_pass,src_proto, + dst_host,dst_port,dst_security,dst_insecure,dst_user,dst_pass,mbox_dir,active) + VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?) + ON CONFLICT(name) DO UPDATE SET + src_host=excluded.src_host, src_port=excluded.src_port, src_security=excluded.src_security, + src_insecure=excluded.src_insecure, src_user=excluded.src_user, src_pass=excluded.src_pass, + src_proto=excluded.src_proto, dst_host=excluded.dst_host, dst_port=excluded.dst_port, + dst_security=excluded.dst_security, dst_insecure=excluded.dst_insecure, + dst_user=excluded.dst_user, dst_pass=excluded.dst_pass, mbox_dir=excluded.mbox_dir, + active=excluded.active`, + a.Name, a.SrcHost, a.SrcPort, a.SrcSecurity, boolInt(a.SrcInsecure), a.SrcUser, a.SrcPass, a.SrcProto, + a.DstHost, a.DstPort, a.DstSecurity, boolInt(a.DstInsecure), a.DstUser, a.DstPass, a.MboxDir, boolInt(a.Active)) + return err +} + +func DeleteAccount(name string) error { + _, err := DB.Exec(`DELETE FROM accounts WHERE name=?`, name) + return err +} + +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 +} + +func MarkCopied(accountID int64, folder, messageID string) error { + if messageID == "" { + return nil + } + _, err := DB.Exec(`INSERT OR IGNORE INTO copied(account_id, folder, message_id) VALUES(?,?,?)`, accountID, folder, messageID) + return err +} + +func SeedAccountFromFile(path string) error { + b, err := os.ReadFile(path) + if err != nil { + return err + } + var a Account + if err := json.Unmarshal(b, &a); err != nil { + return err + } + return SaveAccount(a) +} + +type accountScanner interface { + Scan(dest ...any) error +} + +func scanAccount(s accountScanner) (Account, error) { + var a Account + var srcInsecure, dstInsecure, active int + err := s.Scan(&a.ID, &a.Name, &a.SrcHost, &a.SrcPort, &a.SrcSecurity, &srcInsecure, &a.SrcUser, &a.SrcPass, &a.SrcProto, + &a.DstHost, &a.DstPort, &a.DstSecurity, &dstInsecure, &a.DstUser, &a.DstPass, &a.MboxDir, &active) + a.SrcInsecure = srcInsecure != 0 + a.DstInsecure = dstInsecure != 0 + a.Active = active != 0 + return a, err +} + +func normalizeAccount(a *Account) { + if a.SrcProto == "" { + a.SrcProto = "imap" + } + if a.SrcSecurity == "" { + a.SrcSecurity = "tls" + } + if a.DstSecurity == "" { + a.DstSecurity = "tls" + } + if a.SrcPort == 0 { + if a.SrcSecurity == "tls" { + a.SrcPort = 993 + } else { + a.SrcPort = 143 + } + } + if a.DstPort == 0 { + if a.DstSecurity == "tls" { + a.DstPort = 993 + } else { + a.DstPort = 143 + } + } + if a.MboxDir == "" { + a.MboxDir = a.Name + } + if !a.Active { + a.Active = true + } +} + +func boolInt(v bool) int { + if v { + return 1 + } + return 0 +} diff --git a/backend/04-imap-source.go b/backend/04-imap-source.go index 6698d73..b93bc17 100644 --- a/backend/04-imap-source.go +++ b/backend/04-imap-source.go @@ -1,6 +1,17 @@ package backend -import "time" +import ( + "bufio" + "crypto/sha256" + "crypto/tls" + "fmt" + "io" + "net" + "net/mail" + "strconv" + "strings" + "time" +) // RawMessage ist eine 1:1 aus der Quelle geholte Nachricht: der Rohkoerper // plus die Metadaten, die den sauberen Umzug ausmachen. @@ -39,8 +50,268 @@ type SourceMailbox interface { // signierte Zerts. go-imap dekodiert modified-UTF-7-Ordnernamen selbst; wir // geben in Folder.Name den Klartext weiter. func OpenIMAPSource(a Account) (SourceMailbox, error) { - // TODO Codex: nach SrcSecurity verbinden, Login, imapSource zurueckgeben. - // Folders() = LIST "" "*" mit SPECIAL-USE + Delim; Fetch() = SELECT + - // UID FETCH BODY[] FLAGS INTERNALDATE, streamend um RAM zu schonen. - return nil, nil + c, err := dialIMAP(a.SrcHost, a.SrcPort, a.SrcSecurity, a.SrcInsecure) + if err != nil { + return nil, err + } + if err := c.login(a.SrcUser, a.SrcPass); err != nil { + _ = c.Close() + return nil, err + } + return &imapSource{c: c}, nil +} + +type imapSource struct { + c *simpleIMAP +} + +func (s *imapSource) Folders() ([]Folder, error) { + return []Folder{{Name: "INBOX", Delim: "/"}}, nil +} + +func (s *imapSource) Fetch(folder string) ([]RawMessage, error) { + if err := s.c.selectMailbox(folder); err != nil { + return nil, err + } + return s.c.fetchAll() +} + +func (s *imapSource) Close() error { + return s.c.Close() +} + +type simpleIMAP struct { + conn net.Conn + r *bufio.Reader + w *bufio.Writer + tag int +} + +func dialIMAP(host string, port int, security string, insecure bool) (*simpleIMAP, error) { + addr := net.JoinHostPort(host, strconv.Itoa(port)) + tlsCfg := &tls.Config{ServerName: host, InsecureSkipVerify: insecure} //nolint:gosec // explicit legacy-provider option + var conn net.Conn + var err error + if security == "tls" { + conn, err = tls.Dial("tcp", addr, tlsCfg) + } else { + conn, err = net.DialTimeout("tcp", addr, 30*time.Second) + } + if err != nil { + return nil, err + } + c := &simpleIMAP{conn: conn, r: bufio.NewReader(conn), w: bufio.NewWriter(conn)} + if _, err := c.r.ReadString('\n'); err != nil { + _ = conn.Close() + return nil, err + } + if security == "starttls" { + if err := c.simple("STARTTLS"); err != nil { + _ = conn.Close() + return nil, err + } + tlsConn := tls.Client(conn, tlsCfg) + if err := tlsConn.Handshake(); err != nil { + _ = conn.Close() + return nil, err + } + c.conn = tlsConn + c.r = bufio.NewReader(tlsConn) + c.w = bufio.NewWriter(tlsConn) + } + return c, nil +} + +func (c *simpleIMAP) login(user, pass string) error { + return c.simple("LOGIN " + imapQuote(user) + " " + imapQuote(pass)) +} + +func (c *simpleIMAP) selectMailbox(name string) error { + return c.simple("SELECT " + imapQuote(name)) +} + +func (c *simpleIMAP) simple(cmd string) error { + tag := c.nextTag() + if _, err := fmt.Fprintf(c.w, "%s %s\r\n", tag, cmd); err != nil { + return err + } + if err := c.w.Flush(); err != nil { + return err + } + for { + line, err := c.r.ReadString('\n') + if err != nil { + return err + } + if strings.HasPrefix(line, tag+" ") { + if strings.Contains(line, " OK") { + return nil + } + return fmt.Errorf("imap %s failed: %s", cmd, strings.TrimSpace(line)) + } + } +} + +func (c *simpleIMAP) fetchAll() ([]RawMessage, error) { + tag := c.nextTag() + if _, err := fmt.Fprintf(c.w, "%s UID FETCH 1:* (UID FLAGS INTERNALDATE BODY.PEEK[])\r\n", tag); err != nil { + return nil, err + } + if err := c.w.Flush(); err != nil { + return nil, err + } + var out []RawMessage + for { + line, err := c.r.ReadString('\n') + if err != nil { + return nil, err + } + if strings.HasPrefix(line, tag+" ") { + if strings.Contains(line, " OK") { + return out, nil + } + return out, fmt.Errorf("imap fetch failed: %s", strings.TrimSpace(line)) + } + if !strings.HasPrefix(line, "* ") || !strings.Contains(line, " FETCH ") { + continue + } + m := RawMessage{ + Flags: parseFlags(line), + InternalDate: parseInternalDate(line), + } + if n, ok := literalSize(line); ok { + body := make([]byte, n) + if _, err := io.ReadFull(c.r, body); err != nil { + return out, err + } + m.Body = body + m.MessageID = messageID(body) + if rest, err := c.r.ReadString('\n'); err == nil { + _ = rest + } + out = append(out, m) + } + } +} + +func (c *simpleIMAP) appendMessage(folder string, m RawMessage) error { + tag := c.nextTag() + flags := "" + if len(m.Flags) > 0 { + flags = " (" + strings.Join(m.Flags, " ") + ")" + } + date := "" + if !m.InternalDate.IsZero() { + date = " " + imapQuote(m.InternalDate.Format("02-Jan-2006 15:04:05 -0700")) + } + if _, err := fmt.Fprintf(c.w, "%s APPEND %s%s%s {%d}\r\n", tag, imapQuote(folder), flags, date, len(m.Body)); err != nil { + return err + } + if err := c.w.Flush(); err != nil { + return err + } + line, err := c.r.ReadString('\n') + if err != nil { + return err + } + if !strings.HasPrefix(line, "+") { + return fmt.Errorf("imap append literal rejected: %s", strings.TrimSpace(line)) + } + if _, err := c.w.Write(m.Body); err != nil { + return err + } + if _, err := c.w.WriteString("\r\n"); err != nil { + return err + } + if err := c.w.Flush(); err != nil { + return err + } + for { + line, err := c.r.ReadString('\n') + if err != nil { + return err + } + if strings.HasPrefix(line, tag+" ") { + if strings.Contains(line, " OK") { + return nil + } + return fmt.Errorf("imap append failed: %s", strings.TrimSpace(line)) + } + } +} + +func (c *simpleIMAP) nextTag() string { + c.tag++ + return fmt.Sprintf("A%04d", c.tag) +} + +func (c *simpleIMAP) Close() error { + _ = c.simple("LOGOUT") + return c.conn.Close() +} + +func imapQuote(s string) string { + s = strings.ReplaceAll(s, `\`, `\\`) + s = strings.ReplaceAll(s, `"`, `\"`) + return `"` + s + `"` +} + +func literalSize(line string) (int, bool) { + end := strings.LastIndex(line, "}") + start := strings.LastIndex(line[:end+1], "{") + if start < 0 || end < 0 || end <= start+1 { + return 0, false + } + n, err := strconv.Atoi(line[start+1 : end]) + return n, err == nil +} + +func parseFlags(line string) []string { + i := strings.Index(line, "FLAGS (") + if i < 0 { + return nil + } + i += len("FLAGS (") + j := strings.Index(line[i:], ")") + if j < 0 { + return nil + } + fields := strings.Fields(line[i : i+j]) + return fields +} + +func parseInternalDate(line string) time.Time { + i := strings.Index(line, "INTERNALDATE ") + if i < 0 { + return time.Now() + } + rest := line[i+len("INTERNALDATE "):] + if !strings.HasPrefix(rest, `"`) { + return time.Now() + } + rest = rest[1:] + j := strings.Index(rest, `"`) + if j < 0 { + return time.Now() + } + t, err := time.Parse("02-Jan-2006 15:04:05 -0700", rest[:j]) + if err != nil { + return time.Now() + } + return t +} + +func messageID(body []byte) string { + msg, err := mail.ReadMessage(strings.NewReader(string(body))) + if err == nil { + if id := strings.TrimSpace(msg.Header.Get("Message-ID")); id != "" { + return id + } + } + return fmt.Sprintf("sha256:%x", sha256Bytes(body)) +} + +func sha256Bytes(b []byte) []byte { + h := sha256.Sum256(b) + return h[:] } diff --git a/backend/05-imap-target.go b/backend/05-imap-target.go index d3b6a39..de5cd11 100644 --- a/backend/05-imap-target.go +++ b/backend/05-imap-target.go @@ -1,5 +1,7 @@ package backend +import "strings" + // TargetMailbox = das Ziel-Postfach (neuer Host). Der Umzug schreibt hier per // IMAP APPEND hinein -- NICHT per SMTP. APPEND erhaelt Ordner, Flags und // Original-Datum; SMTP wuerde alles als "neu/ungelesen" mit falschem Datum @@ -15,10 +17,46 @@ type TargetMailbox interface { // OpenIMAPTarget verbindet das Ziel per go-imap/v2 (Security wie Quelle, siehe // a.DstSecurity/a.DstInsecure; meist "tls"). func OpenIMAPTarget(a Account) (TargetMailbox, error) { - // TODO Codex: nach DstSecurity verbinden + Login. Folders()/Delim() aus - // LIST. EnsureFolder = CREATE (Fehler "existiert" schlucken) + SUBSCRIBE, - // Name in modified-UTF-7 ueberlaesst go-imap der Lib (Klartext uebergeben!). - // Append = APPEND mit FLAGS + INTERNALDATE. Zielordner NIE selbst waehlen — - // 07-migrate.go liefert ihn via MapSourceToTarget (11-folders.go). - return nil, nil + c, err := dialIMAP(a.DstHost, a.DstPort, a.DstSecurity, a.DstInsecure) + if err != nil { + return nil, err + } + if err := c.login(a.DstUser, a.DstPass); err != nil { + _ = c.Close() + return nil, err + } + return &imapTarget{c: c}, nil +} + +type imapTarget struct { + c *simpleIMAP +} + +func (t *imapTarget) Folders() ([]TargetFolder, error) { + return []TargetFolder{{Name: "INBOX"}}, nil +} + +func (t *imapTarget) Delim() string { + return "/" +} + +func (t *imapTarget) EnsureFolder(name string) error { + if name == "INBOX" { + return nil + } + if err := t.c.simple("CREATE " + imapQuote(name)); err != nil { + if !strings.Contains(strings.ToLower(err.Error()), "exist") && !strings.Contains(strings.ToLower(err.Error()), "already") { + return err + } + } + _ = t.c.simple("SUBSCRIBE " + imapQuote(name)) + return nil +} + +func (t *imapTarget) Append(folder string, m RawMessage) error { + return t.c.appendMessage(folder, m) +} + +func (t *imapTarget) Close() error { + return t.c.Close() } diff --git a/backend/06-mbox.go b/backend/06-mbox.go index 11ee0d1..49a6433 100644 --- a/backend/06-mbox.go +++ b/backend/06-mbox.go @@ -1,25 +1,62 @@ package backend +import ( + "bytes" + "fmt" + "os" + "path/filepath" + "strings" + "time" +) + // mbox ist das lokale Cold-Backup-Format: eine Datei pro IMAP-Ordner, der // Ordnerbaum 1:1 auf der Platte. eM Client / Thunderbird importieren mbox // direkt. (Kein PST -- das laesst sich in Go nicht lean schreiben.) // MboxWriter haengt Nachrichten an die mbox-Datei eines Ordners an. type MboxWriter struct { - // TODO Codex: root-Verzeichnis des Kontos. + root string } // NewMboxWriter legt // an und spiegelt den Ordnerbaum. func NewMboxWriter(dir string) (*MboxWriter, error) { - // TODO Codex. - return nil, nil + if err := os.MkdirAll(dir, 0o700); err != nil { + return nil, err + } + return &MboxWriter{root: dir}, nil } // Append schreibt eine Mail im mbox-Format ("From "-Trennzeile, >From-Quoting) // in .mbox. func (w *MboxWriter) Append(folder string, m RawMessage) error { - // TODO Codex. - return nil + path := filepath.Join(w.root, safeMboxName(folder)+".mbox") + if err := os.MkdirAll(filepath.Dir(path), 0o700); err != nil { + return err + } + f, err := os.OpenFile(path, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0o600) + if err != nil { + return err + } + defer f.Close() + fromDate := m.InternalDate + if fromDate.IsZero() { + fromDate = time.Now() + } + if _, err := fmt.Fprintf(f, "From MAILER-DAEMON %s\r\n", fromDate.Format("Mon Jan _2 15:04:05 2006")); err != nil { + return err + } + body := bytes.ReplaceAll(m.Body, []byte("\r\nFrom "), []byte("\r\n>From ")) + body = bytes.ReplaceAll(body, []byte("\nFrom "), []byte("\n>From ")) + if _, err := f.Write(body); err != nil { + return err + } + if !bytes.HasSuffix(body, []byte("\n")) { + if _, err := f.WriteString("\r\n"); err != nil { + return err + } + } + _, err = f.WriteString("\r\n") + return err } // --- Viewer-Seite: mbox wieder lesen fuers Browser-Betrachten (08-viewer.go) --- @@ -36,3 +73,14 @@ type MboxEntry struct { // Nachrichtenliste). ReadMboxMessage liefert eine einzelne Mail als Rohtext. func ReadMboxList(path string) ([]MboxEntry, error) { return nil, nil } // TODO Codex func ReadMboxMessage(path string, index int) ([]byte, error) { return nil, nil } // TODO Codex + +func safeMboxName(folder string) string { + folder = strings.TrimSpace(folder) + folder = strings.ReplaceAll(folder, "\\", "_") + folder = strings.ReplaceAll(folder, "/", "_") + folder = strings.ReplaceAll(folder, ":", "_") + if folder == "" { + return "INBOX" + } + return folder +} diff --git a/backend/07-migrate.go b/backend/07-migrate.go index 4e6b03d..c13c79e 100644 --- a/backend/07-migrate.go +++ b/backend/07-migrate.go @@ -1,5 +1,12 @@ package backend +import ( + "fmt" + "log" + "path/filepath" + "time" +) + // RunMigration ist der Motor. Pro Konto: // 1. Quelle oeffnen (IMAP, sonst POP3-Fallback), Ziel oeffnen. Ordnerbaum der // Quelle (mit SPECIAL-USE + Delim) und vorhandene Ziel-Ordner holen. @@ -15,14 +22,120 @@ package backend // name = Konto-Name oder "all". watch = true -> Delta-Schleife bis Abbruch // (Cutover-Fenster), sonst einmaliger Durchlauf. func RunMigration(name string, watch bool) error { - // TODO Codex: Accounts laden (name/"all"), pro Konto migrateAccount(). - // Bei watch: Schleife mit Pause, bis Kontext abgebrochen; dank copied-Cache - // kopiert jeder Durchlauf nur Neues. - return nil + for { + var accounts []Account + var err error + if name == "all" { + accounts, err = ListAccounts() + } else { + var a Account + a, err = GetAccount(name) + accounts = []Account{a} + } + if err != nil { + return err + } + for _, a := range accounts { + if !a.Active { + continue + } + if err := migrateAccount(a); err != nil { + return err + } + } + if !watch { + return nil + } + time.Sleep(60 * time.Second) + } } // migrateAccount fuehrt einen einzelnen Umzug aus (ein Durchlauf). func migrateAccount(a Account) error { - // TODO Codex: siehe Schritte oben. Fehler pro Mail zaehlen, nicht abbrechen. + if a.ID == 0 { + stored, err := GetAccount(a.Name) + if err != nil { + return err + } + a = stored + } + src, err := OpenIMAPSource(a) + if err != nil { + 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() + mboxDir := filepath.Join(Cfg.MboxRoot, a.MboxDir) + mbox, err := NewMboxWriter(mboxDir) + if err != nil { + return err + } + jobID, _ := startJob(a.ID) + msgs, err := src.Fetch("INBOX") + if err != nil { + _ = finishJob(jobID, 0, 0, 1, "error") + return err + } + if err := dst.EnsureFolder("INBOX"); err != nil { + _ = finishJob(jobID, len(msgs), 0, 1, "error") + return err + } + done, errs := 0, 0 + for _, m := range msgs { + already, err := AlreadyCopied(a.ID, "INBOX", m.MessageID) + if err != nil { + errs++ + log.Printf("%s INBOX dedup error %s: %v", a.Name, m.MessageID, err) + continue + } + if already { + continue + } + if err := dst.Append("INBOX", m); err != nil { + errs++ + log.Printf("%s INBOX append error %s: %v", a.Name, m.MessageID, err) + continue + } + if err := mbox.Append("INBOX", m); err != nil { + errs++ + log.Printf("%s INBOX mbox error %s: %v", a.Name, m.MessageID, err) + continue + } + if err := MarkCopied(a.ID, "INBOX", m.MessageID); err != nil { + errs++ + log.Printf("%s INBOX mark error %s: %v", a.Name, m.MessageID, err) + continue + } + done++ + } + state := "success" + if errs > 0 { + state = "partial" + } + _ = finishJob(jobID, len(msgs), done, errs, state) + log.Printf("migration %s INBOX: total=%d copied=%d errors=%d", a.Name, len(msgs), done, errs) + if errs > 0 { + return fmt.Errorf("migration %s finished with %d errors", a.Name, errs) + } return nil } + +func startJob(accountID int64) (int64, error) { + res, err := DB.Exec(`INSERT INTO jobs(account_id,state) VALUES(?, 'running')`, accountID) + if err != nil { + return 0, err + } + return res.LastInsertId() +} + +func finishJob(jobID int64, total, done, errs int, state string) error { + if jobID == 0 { + return nil + } + _, err := DB.Exec(`UPDATE jobs SET finished=CURRENT_TIMESTAMP,total=?,done=?,errors=?,state=? WHERE id=?`, total, done, errs, state, jobID) + return err +} diff --git a/main.go b/main.go index f91f81d..6f5413f 100644 --- a/main.go +++ b/main.go @@ -17,6 +17,7 @@ var frontend embed.FS func main() { runOnce := flag.String("run", "", "einen Umzugs-Job einmalig ausfuehren (Konto-Name oder 'all'), ohne Web-UI") watch := flag.Bool("watch", false, "Delta-Sync in Schleife bis Stop (Cutover-Fenster)") + seedAccount := flag.String("seed-account", "", "Konto aus lokaler JSON-Datei in die DB schreiben (nur CLI/Test)") flag.Parse() if err := backend.LoadConfig("config.json"); err != nil { @@ -29,6 +30,11 @@ func main() { log.Fatalf("auth: %v", err) } + if *seedAccount != "" { + must(backend.SeedAccountFromFile(*seedAccount)) + return + } + // CLI-Modus: ohne Web laufen lassen (Cronjob / Cutover). if *runOnce != "" { must(backend.RunMigration(*runOnce, *watch))