1package store
2
3import (
4 "bufio"
5 "bytes"
6 "errors"
7 "fmt"
8 "io"
9 "maps"
10 "os"
11 "path/filepath"
12 "slices"
13 "strconv"
14 "strings"
15 "time"
16
17 "github.com/mjl-/mox/mlog"
18)
19
20// MsgSource is implemented by readers for mailbox file formats.
21type MsgSource interface {
22 // Return next message, or io.EOF when there are no more.
23 Next() (*Message, *os.File, string, error)
24 Close() error
25}
26
27// MboxReader reads messages from an mbox file, implementing MsgSource.
28type MboxReader struct {
29 log mlog.Log
30 createTemp func(log mlog.Log, pattern string) (*os.File, error)
31 path string
32 line int
33 r *bufio.Reader
34 prevempty bool
35 nonfirst bool
36 eof bool
37 fromLine string // "From "-line for this message.
38 header bool // Now in header section.
39}
40
41// NewMboxReader initializes a MsgSource from which messages can be read.
42func NewMboxReader(log mlog.Log, createTemp func(log mlog.Log, pattern string) (*os.File, error), filename string, r io.Reader) (*MboxReader, error) {
43 return &MboxReader{
44 log: log,
45 createTemp: createTemp,
46 path: filename,
47 line: 1,
48 r: bufio.NewReader(r),
49 }, nil
50}
51
52// Position returns "<filename>:<lineno>" for the current position.
53func (mr *MboxReader) Position() string {
54 return fmt.Sprintf("%s:%d", mr.path, mr.line)
55}
56
57// Next returns the next message read from the mbox file. The file is a temporary
58// file and must be removed/consumed. The third return value is the position in the
59// file.
60func (mr *MboxReader) Next() (*Message, *os.File, string, error) {
61 if mr.eof {
62 return nil, nil, "", io.EOF
63 }
64
65 from := []byte("From ")
66
67 if !mr.nonfirst {
68 mr.header = true
69 // First read, we're at the beginning of the file.
70 line, err := mr.r.ReadBytes('\n')
71 if err == io.EOF {
72 return nil, nil, "", io.EOF
73 }
74 mr.line++
75
76 if !bytes.HasPrefix(line, from) {
77 return nil, nil, mr.Position(), fmt.Errorf(`first line does not start with "From "`)
78 }
79 mr.nonfirst = true
80 mr.fromLine = strings.TrimSpace(string(line))
81 }
82
83 f, err := mr.createTemp(mr.log, "mboxreader")
84 if err != nil {
85 return nil, nil, mr.Position(), err
86 }
87 defer func() {
88 if f != nil {
89 CloseRemoveTempFile(mr.log, f, "message after mbox read error")
90 }
91 }()
92
93 fromLine := mr.fromLine
94 bf := bufio.NewWriter(f)
95 var flags Flags
96 keywords := map[string]bool{}
97 var size int64
98 for {
99 line, err := mr.r.ReadBytes('\n')
100 if err != nil && err != io.EOF {
101 return nil, nil, mr.Position(), fmt.Errorf("reading from mbox: %v", err)
102 }
103 if len(line) > 0 {
104 mr.line++
105 // We store data with crlf, adjust any imported messages with bare newlines. ../rfc/4155:354
106 if !bytes.HasSuffix(line, []byte("\r\n")) {
107 line = append(line[:len(line)-1], "\r\n"...)
108 }
109
110 if mr.header {
111 // See https://doc.dovecot.org/admin_manual/mailbox_formats/mbox/
112 if bytes.HasPrefix(line, []byte("Status:")) {
113 s := strings.TrimSpace(strings.SplitN(string(line), ":", 2)[1])
114 for _, c := range s {
115 switch c {
116 case 'R':
117 flags.Seen = true
118 }
119 }
120 } else if bytes.HasPrefix(line, []byte("X-Status:")) {
121 s := strings.TrimSpace(strings.SplitN(string(line), ":", 2)[1])
122 for _, c := range s {
123 switch c {
124 case 'A':
125 flags.Answered = true
126 case 'F':
127 flags.Flagged = true
128 case 'T':
129 flags.Draft = true
130 case 'D':
131 flags.Deleted = true
132 }
133 }
134 } else if bytes.HasPrefix(line, []byte("X-Keywords:")) {
135 s := strings.TrimSpace(strings.SplitN(string(line), ":", 2)[1])
136 for t := range strings.SplitSeq(s, ",") {
137 word := strings.ToLower(strings.TrimSpace(t))
138 switch word {
139 case "forwarded", "$forwarded":
140 flags.Forwarded = true
141 case "junk", "$junk":
142 flags.Junk = true
143 case "notjunk", "$notjunk", "nonjunk", "$nonjunk":
144 flags.Notjunk = true
145 case "phishing", "$phishing":
146 flags.Phishing = true
147 case "mdnsent", "$mdnsent":
148 flags.MDNSent = true
149 default:
150 if err := CheckKeyword(word); err == nil {
151 keywords[word] = true
152 }
153 }
154 }
155 }
156 }
157 if bytes.Equal(line, []byte("\r\n")) {
158 mr.header = false
159 }
160
161 // Next mail message starts at bare From word. ../rfc/4155:71
162 if mr.prevempty && bytes.HasPrefix(line, from) {
163 mr.fromLine = strings.TrimSpace(string(line))
164 mr.header = true
165 break
166 }
167 // ../rfc/4155:119
168 if bytes.HasPrefix(line, []byte(">")) && bytes.HasPrefix(bytes.TrimLeft(line, ">"), []byte("From ")) {
169 line = line[1:]
170 }
171 n, err := bf.Write(line)
172 if err != nil {
173 return nil, nil, mr.Position(), fmt.Errorf("writing message to file: %v", err)
174 }
175 size += int64(n)
176 mr.prevempty = bytes.Equal(line, []byte("\r\n"))
177 }
178 if err == io.EOF {
179 mr.eof = true
180 break
181 }
182 }
183 if err := bf.Flush(); err != nil {
184 return nil, nil, mr.Position(), fmt.Errorf("flush: %v", err)
185 }
186
187 m := &Message{Flags: flags, Keywords: slices.Sorted(maps.Keys(keywords)), Size: size}
188
189 if t := strings.SplitN(fromLine, " ", 3); len(t) == 3 {
190 layouts := []string{time.ANSIC, time.UnixDate, time.RubyDate}
191 for _, l := range layouts {
192 t, err := time.Parse(l, t[2])
193 if err == nil {
194 m.Received = t
195 break
196 }
197 }
198 }
199
200 // Prevent cleanup by defer.
201 mf := f
202 f = nil
203
204 return m, mf, mr.Position(), nil
205}
206
207// Close is currently a no op, for interface MsgSource.
208func (mr *MboxReader) Close() error {
209 return nil
210}
211
212// we make a slice of files, for cur & new, for sorting by time, so we import
213// messages in a natural order, with most recent messages latest.
214type maildirFile struct {
215 Name string
216 Time time.Time
217}
218
219type MaildirReader struct {
220 log mlog.Log
221 createTemp func(log mlog.Log, pattern string) (*os.File, error)
222 dirNameCur, dirNameNew string
223 rootCur, rootNew *os.Root // For opening files. Closed when Close is called.
224 filesCur, filesNew []maildirFile
225 dovecotFlags []string // Lower-case flags/keywords.
226}
227
228// NewMaildirReader opens the "cur" and "new" files in dir, and returns a MsgSource
229// to read messages from.
230func NewMaildirReader(log mlog.Log, createTemp func(log mlog.Log, pattern string) (*os.File, error), dir string) (*MaildirReader, error) {
231 pathCur := filepath.Join(dir, "cur")
232 pathNew := filepath.Join(dir, "new")
233
234 var rootCur, rootNew *os.Root
235
236 defer func() {
237 if rootCur != nil {
238 err := rootCur.Close()
239 log.Check(err, "closing root for cur dir")
240 }
241 if rootNew != nil {
242 err := rootNew.Close()
243 log.Check(err, "closing root for new dir")
244 }
245 }()
246
247 var err error
248 rootCur, err = os.OpenRoot(pathCur)
249 if err != nil {
250 return nil, fmt.Errorf("open 'cur' path: %w", err)
251 }
252 rootNew, err = os.OpenRoot(pathNew)
253 if err != nil {
254 return nil, fmt.Errorf("open 'new' path: %w", err)
255 }
256
257 filesCur, err := maildirRead(log, pathCur)
258 if err != nil {
259 return nil, fmt.Errorf("reading 'cur' directory: %w", err)
260 }
261 filesNew, err := maildirRead(log, pathNew)
262 if err != nil {
263 return nil, fmt.Errorf("reading 'new' directory: %w", err)
264 }
265
266 mr := &MaildirReader{
267 log: log,
268 createTemp: createTemp,
269 dirNameCur: pathCur,
270 dirNameNew: pathNew,
271 rootCur: rootCur,
272 rootNew: rootNew,
273 filesCur: filesCur,
274 filesNew: filesNew,
275 }
276
277 // Best-effort parsing of dovecot keywords.
278 kf, err := os.Open(filepath.Join(dir, "dovecot-keywords"))
279 if err == nil {
280 mr.dovecotFlags, err = ParseDovecotKeywordsFlags(kf, log)
281 log.Check(err, "parsing dovecot keywords file")
282 err = kf.Close()
283 log.Check(err, "closing dovecot-keywords file")
284 }
285
286 // Prevent cleanup, no more chance of error.
287 rootCur = nil
288 rootNew = nil
289
290 return mr, nil
291}
292
293func maildirRead(log mlog.Log, p string) ([]maildirFile, error) {
294 dir, err := os.Open(p)
295 if err != nil {
296 return nil, err
297 }
298 defer func() {
299 err := dir.Close()
300 log.Check(err, "closing maildir dir")
301 }()
302
303 var files []maildirFile
304 for {
305 ents, err := dir.ReadDir(100)
306 for _, e := range ents {
307 f := maildirFile{
308 Name: e.Name(),
309 Time: messageTime(e),
310 }
311 files = append(files, f)
312 }
313 if err == io.EOF {
314 break
315 } else if err != nil {
316 return nil, fmt.Errorf("read dir: %w", err)
317 }
318 }
319
320 slices.SortFunc(files, func(a, b maildirFile) int { return a.Time.Compare(b.Time) })
321 return files, nil
322}
323
324// Take received time from filename, falling back to mtime for maildirs
325// reconstructed some other sources of message files.
326func messageTime(f os.DirEntry) time.Time {
327 var t time.Time
328 parts := strings.SplitN(f.Name(), ".", 3)
329 if v, err := strconv.ParseInt(parts[0], 10, 64); len(parts) == 3 && err == nil {
330 t = time.Unix(v, 0)
331 } else if fi, err := f.Info(); err == nil {
332 t = fi.ModTime()
333 }
334 return t
335}
336
337func (mr *MaildirReader) Next() (*Message, *os.File, string, error) {
338 var file maildirFile
339 var root *os.Root
340 var dirName string
341 if len(mr.filesCur) > 0 {
342 file = mr.filesCur[0]
343 mr.filesCur = mr.filesCur[1:]
344 root = mr.rootCur
345 dirName = mr.dirNameCur
346 } else if len(mr.filesNew) > 0 {
347 file = mr.filesNew[0]
348 mr.filesNew = mr.filesNew[1:]
349 root = mr.rootNew
350 dirName = mr.dirNameNew
351 } else {
352 return nil, nil, "", io.EOF
353 }
354
355 p := filepath.Join(dirName, file.Name)
356 sf, err := root.Open(file.Name)
357 if err != nil {
358 return nil, nil, p, fmt.Errorf("open message in maildir: %s", err)
359 }
360 defer func() {
361 err := sf.Close()
362 mr.log.Check(err, "closing message file after error")
363 }()
364 f, err := mr.createTemp(mr.log, "maildirreader")
365 if err != nil {
366 return nil, nil, p, err
367 }
368 defer func() {
369 if f != nil {
370 CloseRemoveTempFile(mr.log, f, "maildir temp message file")
371 }
372 }()
373
374 // Copy data, changing bare \n into \r\n.
375 r := bufio.NewReader(sf)
376 w := bufio.NewWriter(f)
377 var size int64
378 for {
379 line, err := r.ReadBytes('\n')
380 if err != nil && err != io.EOF {
381 return nil, nil, p, fmt.Errorf("reading message: %v", err)
382 }
383 if len(line) > 0 {
384 if !bytes.HasSuffix(line, []byte("\r\n")) {
385 line = append(line[:len(line)-1], "\r\n"...)
386 }
387
388 if n, err := w.Write(line); err != nil {
389 return nil, nil, p, fmt.Errorf("writing message: %v", err)
390 } else {
391 size += int64(n)
392 }
393 }
394 if err == io.EOF {
395 break
396 }
397 }
398 if err := w.Flush(); err != nil {
399 return nil, nil, p, fmt.Errorf("writing message: %v", err)
400 }
401
402 // Parse flags. See https://cr.yp.to/proto/maildir.html.
403 flags := Flags{}
404 keywords := map[string]bool{}
405 t := strings.SplitN(file.Name, ":2,", 2)
406 if len(t) == 2 {
407 for _, c := range t[1] {
408 switch c {
409 case 'P':
410 // Passed, doesn't map to a common IMAP flag.
411 case 'R':
412 flags.Answered = true
413 case 'S':
414 flags.Seen = true
415 case 'T':
416 flags.Deleted = true
417 case 'D':
418 flags.Draft = true
419 case 'F':
420 flags.Flagged = true
421 default:
422 if c >= 'a' && c <= 'z' {
423 index := int(c - 'a')
424 if index >= len(mr.dovecotFlags) {
425 continue
426 }
427 kw := mr.dovecotFlags[index]
428 switch kw {
429 case "$forwarded", "forwarded":
430 flags.Forwarded = true
431 case "$junk", "junk":
432 flags.Junk = true
433 case "$notjunk", "notjunk", "nonjunk":
434 flags.Notjunk = true
435 case "$mdnsent", "mdnsent":
436 flags.MDNSent = true
437 case "$phishing", "phishing":
438 flags.Phishing = true
439 default:
440 keywords[kw] = true
441 }
442 }
443 }
444 }
445 }
446
447 m := &Message{Received: file.Time, Flags: flags, Keywords: slices.Sorted(maps.Keys(keywords)), Size: size}
448
449 // Prevent cleanup by defer.
450 mf := f
451 f = nil
452
453 return m, mf, p, nil
454}
455
456// Close closes internal state. It does not close dirNew and dirCur passed to
457// NewMaildirReader.
458func (mr *MaildirReader) Close() error {
459 var err0, err1 error
460 if mr.rootCur != nil {
461 err0 = mr.rootCur.Close()
462 }
463 if mr.rootNew != nil {
464 err1 = mr.rootNew.Close()
465 }
466 return errors.Join(err0, err1)
467}
468
469// ParseDovecotKeywordsFlags attempts to parse a dovecot-keywords file. It only
470// returns valid flags/keywords, as lower-case. If an error is encountered and
471// returned, any keywords that were found are still returned. The returned list has
472// both system/well-known flags and custom keywords.
473func ParseDovecotKeywordsFlags(r io.Reader, log mlog.Log) ([]string, error) {
474 /*
475 If the dovecot-keywords file is present, we parse its additional flags, see
476 https://doc.dovecot.org/admin_manual/mailbox_formats/maildir/
477
478 0 Old
479 1 Junk
480 2 NonJunk
481 3 $Forwarded
482 4 $Junk
483 */
484 keywords := make([]string, 26)
485 end := 0
486 scanner := bufio.NewScanner(r)
487 var errs []string
488 for scanner.Scan() {
489 s := scanner.Text()
490 t := strings.SplitN(s, " ", 2)
491 if len(t) != 2 {
492 errs = append(errs, fmt.Sprintf("unexpected dovecot keyword line: %q", s))
493 continue
494 }
495 v, err := strconv.ParseInt(t[0], 10, 32)
496 if err != nil {
497 errs = append(errs, fmt.Sprintf("unexpected dovecot keyword index: %q", s))
498 continue
499 }
500 if v < 0 || v >= int64(len(keywords)) {
501 errs = append(errs, fmt.Sprintf("dovecot keyword index too big: %q", s))
502 continue
503 }
504 index := int(v)
505 if keywords[index] != "" {
506 errs = append(errs, fmt.Sprintf("duplicate dovecot keyword: %q", s))
507 continue
508 }
509 kw := strings.ToLower(t[1])
510 if !systemWellKnownFlags[kw] {
511 if err := CheckKeyword(kw); err != nil {
512 errs = append(errs, fmt.Sprintf("invalid keyword %q", kw))
513 continue
514 }
515 }
516 keywords[index] = kw
517 if index >= end {
518 end = index + 1
519 }
520 }
521 if err := scanner.Err(); err != nil {
522 errs = append(errs, fmt.Sprintf("reading dovecot keywords file: %v", err))
523 }
524 var err error
525 if len(errs) > 0 {
526 err = errors.New(strings.Join(errs, "; "))
527 }
528 return keywords[:end], err
529}
530