internal/push/push.go
85 lines · 2124 bytes
1package push
2
3import (
4 "context"
5 "log/slog"
6 "time"
7
8 "gitbay.org/gitbay/internal/config"
9 "gitbay.org/gitbay/internal/store"
10)
11
12// DefaultMaxAttempts matches the mailer's: a flaky APNs delays a
13// notification rather than losing it, up to a point.
14const DefaultMaxAttempts = 5
15
16type Deliverer struct {
17 St *store.Store
18 Cl *Client
19 RetryBase time.Duration
20 MaxAttempts int
21}
22
23func New(st *store.Store, cfg config.Push, siteURL string, retryBase time.Duration) (*Deliverer, error) {
24 cl, err := NewClient(cfg, siteURL)
25 if err != nil {
26 return nil, err
27 }
28 return &Deliverer{St: st, Cl: cl, RetryBase: retryBase, MaxAttempts: DefaultMaxAttempts}, nil
29}
30
31// Run drains the push queue until ctx is done.
32func (d *Deliverer) Run(ctx context.Context) {
33 tick := time.NewTicker(2 * time.Second)
34 defer tick.Stop()
35 for {
36 select {
37 case <-ctx.Done():
38 return
39 case <-tick.C:
40 d.drain(ctx)
41 }
42 }
43}
44
45func (d *Deliverer) drain(ctx context.Context) {
46 due, err := d.St.DuePush(20)
47 if err != nil {
48 slog.Error("push: listing due", "err", err)
49 return
50 }
51 for _, q := range due {
52 res, after, sendErr := d.Cl.Send(ctx, q.Token, q.Username, q.Badge, q.Title, q.Body, q.Path)
53 msg := ""
54 if sendErr != nil {
55 msg = sendErr.Error()
56 }
57 switch res {
58 case resultSent:
59 d.St.MarkPushSent(q.ID)
60 case resultReap:
61 // The queued rows cascade with the device.
62 if err := d.St.DeletePushDeviceByToken(q.Token); err != nil {
63 slog.Error("push: reaping device", "device", q.DeviceID, "err", err)
64 }
65 case resultRetry:
66 attempt := q.Attempts + 1
67 if attempt >= d.MaxAttempts {
68 d.St.MarkPushFailed(q.ID, msg, nil)
69 // The device id, never the token.
70 slog.Warn("push dead-lettered",
71 "push", q.ID, "device", q.DeviceID, "attempts", attempt, "err", msg)
72 continue
73 }
74 wait := after
75 if wait == 0 {
76 wait = d.RetryBase << (attempt - 1)
77 }
78 next := time.Now().Add(wait)
79 d.St.MarkPushFailed(q.ID, msg, &next)
80 default: // resultDead
81 d.St.MarkPushFailed(q.ID, msg, nil)
82 slog.Warn("push rejected", "push", q.ID, "device", q.DeviceID, "err", msg)
83 }
84 }
85}