1package webaccount
2
3import (
4 "archive/tar"
5 "archive/zip"
6 "bufio"
7 "bytes"
8 "compress/gzip"
9 "context"
10 cryptrand "crypto/rand"
11 "encoding/json"
12 "errors"
13 "fmt"
14 "io"
15 "log/slog"
16 "maps"
17 "os"
18 "path"
19 "runtime/debug"
20 "slices"
21 "strconv"
22 "strings"
23 "time"
24
25 "golang.org/x/text/unicode/norm"
26
27 "github.com/mjl-/bstore"
28
29 "github.com/mjl-/mox/config"
30 "github.com/mjl-/mox/message"
31 "github.com/mjl-/mox/metrics"
32 "github.com/mjl-/mox/mlog"
33 "github.com/mjl-/mox/mox-"
34 "github.com/mjl-/mox/store"
35)
36
37type importListener struct {
38 Token string
39 Events chan importEvent
40 Register chan bool // Whether register is successful.
41}
42
43type importEvent struct {
44 Token string
45 SSEMsg []byte // Full SSE message, including event: ... and data: ... \n\n
46 Event any // nil, importCount, importProblem, importDone, importAborted
47 Cancel func() // For cancelling the context causing abort of the import. Set in first, import-registering, event.
48}
49
50type importAbortRequest struct {
51 Token string
52 Response chan error
53}
54
55var importers = struct {
56 Register chan *importListener
57 Unregister chan *importListener
58 Events chan importEvent
59 Abort chan importAbortRequest
60 Stop chan struct{}
61}{
62 make(chan *importListener, 1),
63 make(chan *importListener, 1),
64 make(chan importEvent),
65 make(chan importAbortRequest),
66 make(chan struct{}),
67}
68
69// ImportManage should be run as a goroutine, it manages imports of mboxes/maildirs, propagating progress over SSE connections.
70func ImportManage() {
71 log := mlog.New("httpimport", nil)
72 defer func() {
73 if x := recover(); x != nil {
74 log.Error("import manage panic", slog.Any("err", x))
75 debug.PrintStack()
76 metrics.PanicInc(metrics.Importmanage)
77 }
78 }()
79
80 type state struct {
81 MailboxCounts map[string]int
82 Problems []string
83 Done *time.Time
84 Aborted *time.Time
85 Listeners map[*importListener]struct{}
86 Cancel func()
87 }
88
89 imports := map[string]state{} // Token to state.
90 for {
91 select {
92 case l := <-importers.Register:
93 // If we have state, send it so the client is up to date.
94 s, ok := imports[l.Token]
95 l.Register <- ok
96 if !ok {
97 break
98 }
99 s.Listeners[l] = struct{}{}
100
101 sendEvent := func(kind string, v any) {
102 buf, err := json.Marshal(v)
103 if err != nil {
104 log.Errorx("marshal event", err, slog.String("kind", kind), slog.Any("event", v))
105 return
106 }
107 ssemsg := fmt.Sprintf("event: %s\ndata: %s\n\n", kind, buf)
108
109 select {
110 case l.Events <- importEvent{kind, []byte(ssemsg), nil, nil}:
111 default:
112 log.Debug("dropped initial import event to slow consumer")
113 }
114 }
115
116 for m, c := range s.MailboxCounts {
117 sendEvent("count", importCount{m, c})
118 }
119 for _, p := range s.Problems {
120 sendEvent("problem", importProblem{p})
121 }
122 if s.Done != nil {
123 sendEvent("done", importDone{})
124 } else if s.Aborted != nil {
125 sendEvent("aborted", importAborted{})
126 }
127
128 case l := <-importers.Unregister:
129 delete(imports[l.Token].Listeners, l)
130
131 case e := <-importers.Events:
132 s, ok := imports[e.Token]
133 if !ok {
134 s = state{
135 MailboxCounts: map[string]int{},
136 Listeners: map[*importListener]struct{}{},
137 Cancel: e.Cancel,
138 }
139 imports[e.Token] = s
140 }
141 for l := range s.Listeners {
142 select {
143 case l.Events <- e:
144 default:
145 log.Debug("dropped import event to slow consumer")
146 }
147 }
148 if e.Event != nil {
149 s := imports[e.Token]
150 switch x := e.Event.(type) {
151 case importCount:
152 s.MailboxCounts[x.Mailbox] = x.Count
153 case importProblem:
154 s.Problems = append(s.Problems, x.Message)
155 case importDone:
156 now := time.Now()
157 s.Done = &now
158 case importAborted:
159 now := time.Now()
160 s.Aborted = &now
161 }
162 imports[e.Token] = s
163 }
164
165 case a := <-importers.Abort:
166 s, ok := imports[a.Token]
167 if !ok {
168 a.Response <- errors.New("import not found")
169 return
170 }
171 if s.Done != nil {
172 a.Response <- errors.New("import already finished")
173 return
174 }
175 s.Cancel()
176 a.Response <- nil
177
178 case <-importers.Stop:
179 return
180 }
181
182 // Cleanup old state.
183 for t, s := range imports {
184 if len(s.Listeners) > 0 {
185 continue
186 }
187 if s.Done != nil && time.Since(*s.Done) > time.Minute || s.Aborted != nil && time.Since(*s.Aborted) > time.Minute {
188 delete(imports, t)
189 }
190 }
191 }
192}
193
194type importCount struct {
195 Mailbox string
196 Count int
197}
198type importProblem struct {
199 Message string
200}
201type importDone struct{}
202type importAborted struct{}
203type importStep struct {
204 Title string
205}
206
207// importStart prepare the import and launches the goroutine to actually import.
208// importStart is responsible for closing f and removing f.
209func importStart(log mlog.Log, accName string, f *os.File, skipMailboxPrefix string) (string, bool, error) {
210 defer func() {
211 if f != nil {
212 store.CloseRemoveTempFile(log, f, "upload for import")
213 }
214 }()
215
216 var buf [16]byte
217 cryptrand.Read(buf[:])
218 token := fmt.Sprintf("%x", buf)
219
220 if _, err := f.Seek(0, 0); err != nil {
221 return "", false, fmt.Errorf("seek to start of file: %v", err)
222 }
223
224 // Recognize file format.
225 var iszip bool
226 magicZip := []byte{0x50, 0x4b, 0x03, 0x04}
227 magicGzip := []byte{0x1f, 0x8b}
228 magic := make([]byte, 4)
229 if _, err := f.ReadAt(magic, 0); err != nil {
230 return "", true, fmt.Errorf("detecting file format: %v", err)
231 }
232 if bytes.Equal(magic, magicZip) {
233 iszip = true
234 } else if !bytes.Equal(magic[:2], magicGzip) {
235 return "", true, fmt.Errorf("file is not a zip or gzip file")
236 }
237
238 var zr *zip.Reader
239 var tr *tar.Reader
240 if iszip {
241 fi, err := f.Stat()
242 if err != nil {
243 return "", false, fmt.Errorf("stat temporary import zip file: %v", err)
244 }
245 zr, err = zip.NewReader(f, fi.Size())
246 if err != nil {
247 return "", true, fmt.Errorf("opening zip file: %v", err)
248 }
249 } else {
250 gzr, err := gzip.NewReader(f)
251 if err != nil {
252 return "", true, fmt.Errorf("gunzip: %v", err)
253 }
254 tr = tar.NewReader(gzr)
255 }
256
257 acc, err := store.OpenAccount(log, accName, false)
258 if err != nil {
259 return "", false, fmt.Errorf("open acount: %v", err)
260 }
261 acc.Lock() // Not using WithWLock because importMessage is responsible for unlocking.
262
263 tx, err := acc.DB.Begin(context.Background(), true)
264 if err != nil {
265 acc.Unlock()
266 xerr := acc.Close()
267 log.Check(xerr, "closing account")
268 return "", false, fmt.Errorf("start transaction: %v", err)
269 }
270
271 // Ensure token is registered before returning, with context that can be canceled.
272 ctx, cancel := context.WithCancel(mox.Shutdown)
273 importers.Events <- importEvent{token, []byte(": keepalive\n\n"), nil, cancel}
274
275 log.Info("starting import")
276 go importMessages(ctx, log.WithCid(mox.Cid()), token, acc, tx, zr, tr, f, skipMailboxPrefix)
277 f = nil // importMessages is now responsible for closing and removing.
278
279 return token, false, nil
280}
281
282// importMessages imports the messages from zip/tgz file f.
283// importMessages is responsible for unlocking and closing acc, and closing tx and f.
284func importMessages(ctx context.Context, log mlog.Log, token string, acc *store.Account, tx *bstore.Tx, zr *zip.Reader, tr *tar.Reader, f *os.File, skipMailboxPrefix string) {
285 // If a fatal processing error occurs, we panic with this type.
286 type importError struct{ Err error }
287
288 // During import we collect all changes and broadcast them at the end, when successful.
289 var changes []store.Change
290
291 // ID's of delivered messages. If we have to rollback, we have to remove this files.
292 var newIDs []int64
293
294 sendEvent := func(kind string, v any) {
295 buf, err := json.Marshal(v)
296 if err != nil {
297 log.Errorx("marshal event", err, slog.String("kind", kind), slog.Any("event", v))
298 return
299 }
300 ssemsg := fmt.Sprintf("event: %s\ndata: %s\n\n", kind, buf)
301 importers.Events <- importEvent{token, []byte(ssemsg), v, nil}
302 }
303
304 canceled := func() bool {
305 select {
306 case <-ctx.Done():
307 sendEvent("aborted", importAborted{})
308 return true
309 default:
310 return false
311 }
312 }
313
314 problemf := func(format string, args ...any) {
315 msg := fmt.Sprintf(format, args...)
316 sendEvent("problem", importProblem{Message: msg})
317 }
318
319 defer func() {
320 store.CloseRemoveTempFile(log, f, "uploaded messages")
321
322 for _, id := range newIDs {
323 p := acc.MessagePath(id)
324 err := os.Remove(p)
325 log.Check(err, "closing message file after import error", slog.String("path", p))
326 }
327 if tx != nil {
328 err := tx.Rollback()
329 log.Check(err, "rolling back transaction")
330 }
331 if acc != nil {
332 acc.Unlock()
333 err := acc.Close()
334 log.Check(err, "closing account")
335 }
336
337 x := recover()
338 if x == nil {
339 return
340 }
341 if err, ok := x.(importError); ok {
342 log.Errorx("import error", err.Err)
343 problemf("%s (aborting)", err.Err)
344 sendEvent("aborted", importAborted{})
345 } else {
346 log.Error("import panic", slog.Any("err", x))
347 debug.PrintStack()
348 metrics.PanicInc(metrics.Importmessages)
349 }
350 }()
351
352 ximportcheckf := func(err error, format string, args ...any) {
353 if err != nil {
354 panic(importError{fmt.Errorf("%s: %s", fmt.Sprintf(format, args...), err)})
355 }
356 }
357
358 err := acc.ThreadingWait(log)
359 ximportcheckf(err, "waiting for account thread upgrade")
360
361 conf, _ := acc.Conf()
362
363 jf, _, err := acc.OpenJunkFilter(ctx, log)
364 if err != nil && !errors.Is(err, store.ErrNoJunkFilter) {
365 ximportcheckf(err, "open junk filter")
366 }
367 defer func() {
368 if jf != nil {
369 err := jf.CloseDiscard()
370 log.Check(err, "closing junk filter")
371 }
372 }()
373
374 // Mailboxes we imported, and message counts.
375 mailboxNames := map[string]*store.Mailbox{}
376 mailboxIDs := map[int64]*store.Mailbox{}
377 mailboxKeywordCounts := map[int64]int{}
378 messages := map[string]int{}
379
380 maxSize := acc.QuotaMessageSize()
381 du := store.DiskUsage{ID: 1}
382 err = tx.Get(&du)
383 ximportcheckf(err, "get disk usage")
384 var addSize int64
385
386 // For maildirs, we are likely to get a possible dovecot-keywords file after having
387 // imported the messages. Once we see the keywords, we use them. But before that
388 // time we remember which messages miss keywords. Once the keywords become
389 // available, we'll fix up the flags for the unknown messages
390 mailboxKeywords := map[string]map[rune]string{} // Mailbox to 'a'-'z' to flag name.
391 mailboxMissingKeywordMessages := map[string]map[int64]string{} // Mailbox to message id to string consisting of the unrecognized flags.
392
393 // Previous mailbox an event was sent for. We send an event for new mailboxes, when
394 // another 100 messages were added, when adding a message to another mailbox, and
395 // finally at the end as a closing statement.
396 var prevMailbox string
397
398 var modseq store.ModSeq // Assigned on first message, used for all messages.
399
400 trainMessage := func(m *store.Message, p message.Part, pos string) {
401 words, err := jf.ParseMessage(p)
402 if err != nil {
403 problemf("parsing message %s for updating junk filter: %v (continuing)", pos, err)
404 return
405 }
406 err = jf.Train(ctx, !m.Junk, words)
407 if err != nil {
408 problemf("training junk filter for message %s: %v (continuing)", pos, err)
409 return
410 }
411 m.TrainedJunk = &m.Junk
412 }
413
414 openTrainMessage := func(m *store.Message) {
415 path := acc.MessagePath(m.ID)
416 f, err := os.Open(path)
417 if err != nil {
418 problemf("opening message again for training junk filter: %v (continuing)", err)
419 return
420 }
421 defer func() {
422 err := f.Close()
423 log.Check(err, "closing file after training junkfilter")
424 }()
425 p, err := m.LoadPart(f)
426 if err != nil {
427 problemf("loading parsed message again for training junk filter: %v (continuing)", err)
428 return
429 }
430 trainMessage(m, p, fmt.Sprintf("message id %d", m.ID))
431 }
432
433 xensureMailbox := func(name string) *store.Mailbox {
434 // Ensure name is normalized.
435 name = norm.NFC.String(name)
436 name, _, err := config.CheckMailboxName(name, true)
437 ximportcheckf(err, "checking mailbox name")
438
439 if mb, ok := mailboxNames[name]; ok {
440 return mb
441 }
442
443 var p string
444 var mb *store.Mailbox
445 var parent store.Mailbox
446 for i, e := range strings.Split(name, "/") {
447 if i == 0 {
448 p = e
449 } else {
450 p = path.Join(p, e)
451 }
452 if _, ok := mailboxNames[p]; ok {
453 continue
454 }
455
456 mb, err = acc.MailboxFind(tx, p)
457 ximportcheckf(err, "looking up mailbox %s to import to (aborting)", p)
458 if mb == nil {
459 uidvalidity, err := acc.NextUIDValidity(tx)
460 ximportcheckf(err, "finding next uid validity")
461
462 if modseq == 0 {
463 var err error
464 modseq, err = acc.NextModSeq(tx)
465 ximportcheckf(err, "assigning next modseq")
466 }
467
468 mb = &store.Mailbox{
469 CreateSeq: modseq,
470 ModSeq: modseq,
471 ParentID: parent.ID,
472 Name: p,
473 UIDValidity: uidvalidity,
474 UIDNext: 1,
475 HaveCounts: true,
476 // Do not assign special-use flags. This existing account probably already has such mailboxes.
477 }
478 err = tx.Insert(mb)
479 ximportcheckf(err, "inserting mailbox in database")
480 parent = *mb
481
482 if tx.Get(&store.Subscription{Name: p}) != nil {
483 err := tx.Insert(&store.Subscription{Name: p})
484 ximportcheckf(err, "subscribing to imported mailbox")
485 }
486 changes = append(changes, store.ChangeAddMailbox{Mailbox: *mb, Flags: []string{`\Subscribed`}})
487 }
488 if prevMailbox != "" && mb.Name != prevMailbox {
489 sendEvent("count", importCount{prevMailbox, messages[prevMailbox]})
490 }
491 mailboxKeywordCounts[mb.ID] = len(mb.Keywords)
492 mailboxNames[mb.Name] = mb
493 mailboxIDs[mb.ID] = mb
494 sendEvent("count", importCount{mb.Name, 0})
495 prevMailbox = mb.Name
496 }
497 return mb
498 }
499
500 xdeliver := func(mb *store.Mailbox, m *store.Message, f *os.File, pos string) {
501 defer store.CloseRemoveTempFile(log, f, "message file for import")
502 m.MailboxID = mb.ID
503 m.MailboxOrigID = mb.ID
504
505 addSize += m.Size
506 if maxSize > 0 && du.MessageSize+addSize > maxSize {
507 ximportcheckf(fmt.Errorf("account over maximum total size %d", maxSize), "checking quota")
508 }
509
510 if modseq == 0 {
511 var err error
512 modseq, err = acc.NextModSeq(tx)
513 ximportcheckf(err, "assigning next modseq")
514 }
515 m.CreateSeq = modseq
516 m.ModSeq = modseq
517
518 // Parse message and store parsed information for later fast retrieval.
519 p, err := message.EnsurePart(log.Logger, false, f, m.Size)
520 if err != nil {
521 problemf("parsing message %s: %s (continuing)", pos, err)
522 }
523 m.ParsedBuf, err = json.Marshal(p)
524 ximportcheckf(err, "marshal parsed message structure")
525
526 // Set fields needed for future threading. By doing it now, MessageAdd won't
527 // have to parse the Part again.
528 p.SetReaderAt(store.FileMsgReader(m.MsgPrefix, f))
529 m.PrepareThreading(log, &p)
530
531 if m.Received.IsZero() {
532 if p.Envelope != nil && !p.Envelope.Date.IsZero() {
533 m.Received = p.Envelope.Date
534 } else {
535 m.Received = time.Now()
536 }
537 }
538
539 // We set the flags that Deliver would set now and train ourselves. This prevents
540 // Deliver from training, which would open the junk filter, change it, and write it
541 // back to disk, for each message (slow).
542 m.JunkFlagsForMailbox(*mb, conf)
543 if jf != nil && m.NeedsTraining() {
544 trainMessage(m, p, pos)
545 }
546
547 opts := store.AddOpts{
548 SkipDirSync: true,
549 SkipTraining: true,
550 SkipThreads: true,
551 SkipUpdateDiskUsage: true,
552 SkipCheckQuota: true,
553 SkipPreview: true,
554 }
555 if err := acc.MessageAdd(log, tx, mb, m, f, opts); err != nil {
556 problemf("delivering message %s: %s (continuing)", pos, err)
557 return
558 }
559 newIDs = append(newIDs, m.ID)
560 changes = append(changes, m.ChangeAddUID(*mb))
561 messages[mb.Name]++
562 if messages[mb.Name]%100 == 0 || prevMailbox != mb.Name {
563 prevMailbox = mb.Name
564 sendEvent("count", importCount{mb.Name, messages[mb.Name]})
565 }
566 }
567
568 ximportMbox := func(mailbox, filename string, r io.Reader) {
569 if mailbox == "" {
570 problemf("empty mailbox name for mbox file %s (skipping)", filename)
571 return
572 }
573 mb := xensureMailbox(mailbox)
574
575 mr := store.NewMboxReader(log, store.CreateMessageTemp, filename, r)
576 for {
577 m, mf, pos, err := mr.Next()
578 if err == io.EOF {
579 break
580 } else if err != nil {
581 ximportcheckf(err, "next message in mbox file")
582 }
583
584 xdeliver(mb, m, mf, pos)
585 }
586 }
587
588 ximportMaildir := func(mailbox, filename string, r io.Reader) {
589 if mailbox == "" {
590 problemf("empty mailbox name for maildir file %s (skipping)", filename)
591 return
592 }
593 mb := xensureMailbox(mailbox)
594
595 f, err := store.CreateMessageTemp(log, "import")
596 ximportcheckf(err, "creating temp message")
597 defer func() {
598 if f != nil {
599 store.CloseRemoveTempFile(log, f, "message to import")
600 }
601 }()
602
603 // Copy data, changing bare \n into \r\n.
604 br := bufio.NewReader(r)
605 w := bufio.NewWriter(f)
606 var size int64
607 for {
608 line, err := br.ReadBytes('\n')
609 if err != nil && err != io.EOF {
610 ximportcheckf(err, "reading message")
611 }
612 if len(line) > 0 {
613 if !bytes.HasSuffix(line, []byte("\r\n")) {
614 line = append(line[:len(line)-1], "\r\n"...)
615 }
616
617 n, err := w.Write(line)
618 ximportcheckf(err, "writing message")
619 size += int64(n)
620 }
621 if err == io.EOF {
622 break
623 }
624 }
625 err = w.Flush()
626 ximportcheckf(err, "writing message")
627
628 var received time.Time
629 t := strings.SplitN(path.Base(filename), ".", 2)
630 if v, err := strconv.ParseInt(t[0], 10, 64); err == nil {
631 received = time.Unix(v, 0)
632 }
633
634 // Parse flags. See https://cr.yp.to/proto/maildir.html.
635 var keepFlags strings.Builder
636 var flags store.Flags
637 keywords := map[string]bool{}
638 t = strings.SplitN(path.Base(filename), ":2,", 2)
639 if len(t) == 2 {
640 for _, c := range t[1] {
641 switch c {
642 case 'P':
643 // Passed, doesn't map to a common IMAP flag.
644 case 'R':
645 flags.Answered = true
646 case 'S':
647 flags.Seen = true
648 case 'T':
649 flags.Deleted = true
650 case 'D':
651 flags.Draft = true
652 case 'F':
653 flags.Flagged = true
654 default:
655 if c >= 'a' && c <= 'z' {
656 dovecotKeywords, ok := mailboxKeywords[mailbox]
657 if !ok {
658 // No keywords file seen yet, we'll try later if it comes in.
659 keepFlags.WriteString(string(c))
660 } else if kw, ok := dovecotKeywords[c]; ok {
661 flagSet(&flags, keywords, kw)
662 }
663 }
664 }
665 }
666 }
667
668 m := store.Message{
669 Received: received,
670 Flags: flags,
671 Keywords: slices.Sorted(maps.Keys(keywords)),
672 Size: size,
673 }
674 xdeliver(mb, &m, f, filename)
675 f = nil
676 if keepFlags.String() != "" {
677 if _, ok := mailboxMissingKeywordMessages[mailbox]; !ok {
678 mailboxMissingKeywordMessages[mailbox] = map[int64]string{}
679 }
680 mailboxMissingKeywordMessages[mailbox][m.ID] = keepFlags.String()
681 }
682 }
683
684 importFile := func(name string, r io.Reader) {
685 origName := name
686
687 if strings.HasPrefix(name, skipMailboxPrefix) {
688 name = strings.TrimPrefix(name[len(skipMailboxPrefix):], "/")
689 }
690
691 if before, ok := strings.CutSuffix(name, "/"); ok {
692 name = before
693 dir := path.Dir(name)
694 switch path.Base(dir) {
695 case "new", "cur", "tmp":
696 // Maildir, ensure it exists.
697 mailbox := path.Dir(dir)
698 xensureMailbox(mailbox)
699 }
700 // Otherwise, this is just a directory that probably holds mbox files and maildirs.
701 return
702 }
703
704 if strings.HasSuffix(path.Base(name), ".mbox") {
705 mailbox := name[:len(name)-len(".mbox")]
706 ximportMbox(mailbox, origName, r)
707 return
708 }
709 dir := path.Dir(name)
710 dirbase := path.Base(dir)
711 switch dirbase {
712 case "new", "cur", "tmp":
713 mailbox := path.Dir(dir)
714 ximportMaildir(mailbox, origName, r)
715 return
716 }
717
718 if path.Base(name) != "dovecot-keywords" {
719 problemf("unrecognized file %s (skipping)", origName)
720 return
721 }
722
723 // Handle dovecot-keywords.
724 mailbox := path.Dir(name)
725 dovecotKeywords := map[rune]string{}
726 words, err := store.ParseDovecotKeywordsFlags(r, log)
727 log.Check(err, "parsing dovecot keywords for mailbox", slog.String("mailbox", mailbox))
728 for i, kw := range words {
729 dovecotKeywords['a'+rune(i)] = kw
730 }
731 mailboxKeywords[mailbox] = dovecotKeywords
732
733 for id, chars := range mailboxMissingKeywordMessages[mailbox] {
734 var flags, zeroflags store.Flags
735 keywords := map[string]bool{}
736 for _, c := range chars {
737 kw, ok := dovecotKeywords[c]
738 if !ok {
739 problemf("unspecified dovecot message flag %c for message id %d (continuing)", c, id)
740 continue
741 }
742 flagSet(&flags, keywords, kw)
743 }
744 if flags == zeroflags && len(keywords) == 0 {
745 continue
746 }
747
748 m := store.Message{ID: id}
749 err := tx.Get(&m)
750 ximportcheckf(err, "get imported message for flag update")
751
752 mb := mailboxIDs[m.MailboxID]
753 mb.Sub(m.MailboxCounts())
754
755 oflags := m.Flags
756 m.Flags = m.Flags.Set(flags, flags)
757 m.Keywords = slices.Sorted(maps.Keys(keywords))
758
759 mb.Add(m.MailboxCounts())
760
761 mb.Keywords, _ = store.MergeKeywords(mb.Keywords, m.Keywords)
762
763 // We train before updating, training may set m.TrainedJunk.
764 if jf != nil && m.NeedsTraining() {
765 openTrainMessage(&m)
766 }
767 err = tx.Update(&m)
768 ximportcheckf(err, "updating message after flag update")
769 changes = append(changes, m.ChangeFlags(oflags, *mb))
770 }
771 delete(mailboxMissingKeywordMessages, mailbox)
772 }
773
774 if zr != nil {
775 for _, f := range zr.File {
776 if canceled() {
777 return
778 }
779 zf, err := f.Open()
780 if err != nil {
781 problemf("opening file %s in zip: %v", f.Name, err)
782 continue
783 }
784 importFile(f.Name, zf)
785 err = zf.Close()
786 log.Check(err, "closing file from zip")
787 }
788 } else {
789 for {
790 if canceled() {
791 return
792 }
793 h, err := tr.Next()
794 if err == io.EOF {
795 break
796 } else if err != nil {
797 problemf("reading next tar header: %v (aborting)", err)
798 return
799 }
800 importFile(h.Name, tr)
801 }
802 }
803
804 total := 0
805 for _, count := range messages {
806 total += count
807 }
808 log.Debug("messages imported", slog.Int("total", total))
809
810 // Send final update for count of last-imported mailbox.
811 if prevMailbox != "" {
812 sendEvent("count", importCount{prevMailbox, messages[prevMailbox]})
813 }
814
815 // Match threads.
816 if len(newIDs) > 0 {
817 sendEvent("step", importStep{"matching messages with threads"})
818 err = acc.AssignThreads(ctx, log, tx, newIDs[0], 0, io.Discard)
819 ximportcheckf(err, "assigning messages to threads")
820 }
821
822 // Update mailboxes with counts and keywords.
823 for _, mb := range mailboxIDs {
824 err = tx.Update(mb)
825 ximportcheckf(err, "updating mailbox count and keywords")
826
827 changes = append(changes, mb.ChangeCounts())
828 if len(mb.Keywords) != mailboxKeywordCounts[mb.ID] {
829 changes = append(changes, mb.ChangeKeywords())
830 }
831 }
832
833 err = acc.AddMessageSize(log, tx, addSize)
834 ximportcheckf(err, "updating disk usage after import")
835
836 err = tx.Commit()
837 tx = nil
838 ximportcheckf(err, "commit")
839 newIDs = nil
840
841 if jf != nil {
842 if err := jf.Close(); err != nil {
843 problemf("saving changes of training junk filter: %v (continuing)", err)
844 log.Errorx("saving changes of training junk filter", err)
845 }
846 jf = nil
847 }
848
849 store.BroadcastChanges(acc, changes)
850 acc.Unlock()
851 err = acc.Close()
852 log.Check(err, "closing account after import")
853 acc = nil
854
855 sendEvent("done", importDone{})
856}
857
858func flagSet(flags *store.Flags, keywords map[string]bool, word string) {
859 switch word {
860 case "forwarded", "$forwarded":
861 flags.Forwarded = true
862 case "junk", "$junk":
863 flags.Junk = true
864 case "notjunk", "$notjunk", "nonjunk", "$nonjunk":
865 flags.Notjunk = true
866 case "phishing", "$phishing":
867 flags.Phishing = true
868 case "mdnsent", "$mdnsent":
869 flags.MDNSent = true
870 default:
871 if err := store.CheckKeyword(word); err == nil {
872 keywords[word] = true
873 }
874 }
875}
876