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}