internal/imap/submission.go (view raw)
1package imap
2
3import (
4 "context"
5 "errors"
6 "fmt"
7 "io"
8 "log/slog"
9 "net"
10 "time"
11
12 "postern/internal/db"
13 "postern/internal/model"
14 "postern/internal/policy"
15)
16
17func (i *Server) SubmitMessage(ctx context.Context, conn net.Conn, user model.User, r io.Reader) error {
18 logger := slog.With("component", "internal_submission", "user_id", user.ID)
19 internalTime := time.Now()
20
21 fileName, size, parsedMsg, err := i.persistence.WriteBlobMessage(r)
22 if err != nil {
23 return fmt.Errorf("write blob: %w", err)
24 }
25
26 if parsedMsg == nil {
27 return errors.New("parsed message is nil")
28 }
29
30 policyUser := policy.UserToStarlark(user, i.db)
31 policyConn := policy.ConnToStarlark(conn)
32
33 fromAddr, err := policy.ParseAddress(parsedMsg.Header.Get("From"))
34 if err != nil {
35 return model.ErrMissingFromHeader
36 }
37 toAddr, err := policy.ParseAddress(parsedMsg.Header.Get("To"))
38 if err != nil {
39 return errors.New("missing To header")
40 }
41 subject := parsedMsg.Header.Get("Subject")
42 msgSize := size
43
44 msg := policy.MessageContext{
45 HeaderFrom: fromAddr,
46 HeaderTo: toAddr,
47 Subject: subject,
48 Headers: parsedMsg.Header,
49 Size: msgSize,
50 }
51
52 destinationMailbox, err := i.policyEngine.OnMessageDeliver(policyConn, policyUser, msg)
53 if destinationMailbox == "" || err != nil {
54 logger.Warn("Policy script error, falling back to INBOX", "error", err)
55 destinationMailbox = "INBOX"
56 }
57
58 logger.InfoContext(ctx, "delivery", "mailbox", destinationMailbox)
59
60 // Create mailbox if not exists
61 mailboxID, _, err := i.db.GetMailboxID(ctx, user.ID, destinationMailbox)
62 if err != nil {
63 if errors.Is(err, db.ErrMailboxNotFound) {
64 err = i.db.CreateMailbox(ctx, user.ID, destinationMailbox)
65 if err != nil {
66 return fmt.Errorf("create on-demand mailbox: %w", err)
67 }
68 mailboxID, _, err = i.db.GetMailboxID(ctx, user.ID, destinationMailbox)
69 if err != nil {
70 return fmt.Errorf("get on-demand mailbox ID: %w", err)
71 }
72 err = i.db.Subscribe(ctx, user.ID, mailboxID, destinationMailbox)
73 if err != nil {
74 return fmt.Errorf("subscribe on-demand mailbox: %w", err)
75 }
76 } else {
77 return fmt.Errorf("get mailbox ID: %w", err)
78 }
79 }
80
81 _, err = i.db.AppendMessage(ctx, mailboxID, user.ID, fileName, size, internalTime.UTC().Format(time.RFC3339), nil, parsedMsg)
82 if err != nil {
83 return fmt.Errorf("append message: %w", err)
84 }
85
86 // Notify user
87 t := i.mTracker.get(user.ID, mailboxID)
88 if t != nil {
89 mb, err := i.db.GetUserMailbox(ctx, user.ID, destinationMailbox)
90 if err != nil {
91 return fmt.Errorf("get user mailbox: %w", err)
92 }
93 t.QueueNumMessages(mb.NumMessages)
94 logger.Info("EXISTS message sent to user", "user", user.ID, "mailbox", mb.Name, "num_messages", mb.NumMessages)
95 }
96
97 return nil
98}