mirror of
https://github.com/multipleof4/ntfy.git
synced 2026-10-08 21:05:21 +00:00
Add message-poll-limit ...
This commit is contained in:
+5
-1
@@ -4,6 +4,7 @@ import (
|
||||
"crypto/sha256"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"heckel.io/ntfy/v2/message"
|
||||
"io/fs"
|
||||
"net/netip"
|
||||
"reflect"
|
||||
@@ -66,7 +67,8 @@ func banWeight(err *errHTTP, weight int) string {
|
||||
// - total topic limit: max number of topics overall
|
||||
// - various attachment limits
|
||||
const (
|
||||
DefaultMessageSizeLimit = 4096 // Bytes; note that FCM/APNS have a limit of ~4 KB for the entire message
|
||||
DefaultMessageSizeLimit = 4096 // Bytes; note that FCM/APNS have a limit of ~4 KB for the entire message
|
||||
DefaultMessagePollLimit = message.NoLimit // Uncapped by default; operators cap it to bound replay memory (limit x MessageSizeLimit)
|
||||
DefaultTotalTopicLimit = 15000
|
||||
DefaultAttachmentTotalSizeLimit = int64(5 * 1024 * 1024 * 1024) // 5 GB
|
||||
DefaultAttachmentFileSizeLimit = int64(15 * 1024 * 1024) // 15 MB
|
||||
@@ -174,6 +176,7 @@ type Config struct {
|
||||
MessageDelayMin time.Duration
|
||||
MessageDelayMax time.Duration
|
||||
MessageSizeLimit int
|
||||
MessagePollLimit int
|
||||
TotalTopicLimit int
|
||||
TotalAttachmentSizeLimit int64
|
||||
VisitorSubscriptionLimit int
|
||||
@@ -282,6 +285,7 @@ func NewConfig() *Config {
|
||||
TwilioVerifyService: "",
|
||||
TwilioCallFormat: nil,
|
||||
MessageSizeLimit: DefaultMessageSizeLimit,
|
||||
MessagePollLimit: DefaultMessagePollLimit,
|
||||
MessageDelayMin: DefaultMessageDelayMin,
|
||||
MessageDelayMax: DefaultMessageDelayMax,
|
||||
TotalTopicLimit: DefaultTotalTopicLimit,
|
||||
|
||||
+15
-7
@@ -1474,7 +1474,7 @@ func (s *Server) handleSubscribeHTTP(w http.ResponseWriter, r *http.Request, v *
|
||||
t.Keepalive()
|
||||
}
|
||||
meterPollBandwidth = true
|
||||
return s.sendOldMessages(topics, since, scheduled, v, sub)
|
||||
return s.sendOldMessages(w, topics, since, scheduled, v, sub)
|
||||
}
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
@@ -1490,7 +1490,7 @@ func (s *Server) handleSubscribeHTTP(w http.ResponseWriter, r *http.Request, v *
|
||||
if err := sub(v, model.NewOpenMessage(topicsStr)); err != nil { // Send out open message
|
||||
return err
|
||||
}
|
||||
if err := s.sendOldMessages(topics, since, scheduled, v, sub); err != nil {
|
||||
if err := s.sendOldMessages(w, topics, since, scheduled, v, sub); err != nil {
|
||||
return err
|
||||
}
|
||||
for {
|
||||
@@ -1625,7 +1625,7 @@ func (s *Server) handleSubscribeWS(w http.ResponseWriter, r *http.Request, v *vi
|
||||
for _, t := range topics {
|
||||
t.Keepalive()
|
||||
}
|
||||
return s.sendOldMessages(topics, since, scheduled, v, sub)
|
||||
return s.sendOldMessages(w, topics, since, scheduled, v, sub)
|
||||
}
|
||||
subscriberIDs := make([]int, 0)
|
||||
for _, t := range topics {
|
||||
@@ -1639,7 +1639,7 @@ func (s *Server) handleSubscribeWS(w http.ResponseWriter, r *http.Request, v *vi
|
||||
if err := sub(v, model.NewOpenMessage(topicsStr)); err != nil { // Send out open message
|
||||
return err
|
||||
}
|
||||
if err := s.sendOldMessages(topics, since, scheduled, v, sub); err != nil {
|
||||
if err := s.sendOldMessages(w, topics, since, scheduled, v, sub); err != nil {
|
||||
return err
|
||||
}
|
||||
err = g.Wait()
|
||||
@@ -1734,21 +1734,29 @@ func (s *Server) setRateVisitors(r *http.Request, v *visitor, rateTopics []*topi
|
||||
|
||||
// sendOldMessages selects old messages from the messageCache and calls sub for each of them. It uses since as the
|
||||
// marker, returning only messages that are newer than the marker.
|
||||
func (s *Server) sendOldMessages(topics []*topic, since model.SinceMarker, scheduled bool, v *visitor, sub subscriber) error {
|
||||
func (s *Server) sendOldMessages(w http.ResponseWriter, topics []*topic, since model.SinceMarker, scheduled bool, v *visitor, sub subscriber) error {
|
||||
if since.IsNone() {
|
||||
return nil
|
||||
}
|
||||
messages := make([]*model.Message, 0)
|
||||
truncated := false
|
||||
for _, t := range topics {
|
||||
topicMessages, err := s.messageCache.Messages(t.ID, since, scheduled)
|
||||
topicMessages, topicTruncated, err := s.messageCache.MessagesCapped(t.ID, since, scheduled, s.config.MessagePollLimit)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
truncated = truncated || topicTruncated
|
||||
messages = append(messages, topicMessages...)
|
||||
}
|
||||
sort.Slice(messages, func(i, j int) bool {
|
||||
// Stable: Time has second granularity, so a multi-topic replay has many equal keys. An unstable
|
||||
// sort reorders them and a topic's own messages come back out of publish order (#1297).
|
||||
sort.SliceStable(messages, func(i, j int) bool {
|
||||
return messages[i].Time < messages[j].Time
|
||||
})
|
||||
// Must be set before the first message is written, or the header is already on the wire
|
||||
if truncated && w != nil {
|
||||
w.Header().Set("X-Messages-Truncated", "1")
|
||||
}
|
||||
for _, m := range messages {
|
||||
if err := sub(v, m); err != nil {
|
||||
return err
|
||||
|
||||
@@ -327,6 +327,15 @@
|
||||
# - message-delay-limit defines the max delay of a message when using the "Delay" header.
|
||||
#
|
||||
# message-size-limit: "4k"
|
||||
|
||||
# Max number of cached messages returned per topic when the cache is replayed, i.e. a poll or a
|
||||
# "since" request. A poll without a "since" cursor returns a topic's entire cache, which is
|
||||
# unbounded in size. Uncapped by default, to preserve the historical behavior; set it to bound
|
||||
# how large a single replay can get. The newest messages are kept, and a truncated response
|
||||
# carries an "X-Messages-Truncated: 1" header. Together with message-size-limit this bounds the
|
||||
# memory one replay can allocate: limit x message-size-limit (e.g. 5000 x 4096 = ~20 MB).
|
||||
#
|
||||
# message-poll-limit: 5000
|
||||
# message-delay-limit: "3d"
|
||||
|
||||
# Rate limiting: Total number of topics before the server rejects new topics.
|
||||
|
||||
@@ -16,6 +16,7 @@ import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
"runtime/debug"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
@@ -2810,6 +2811,77 @@ func TestServer_PublishAttachmentBandwidthLimit(t *testing.T) {
|
||||
})
|
||||
}
|
||||
|
||||
func TestServer_PollOrderAcrossTopics(t *testing.T) {
|
||||
forEachBackend(t, func(t *testing.T, databaseURL string) {
|
||||
// Replaying several topics at once concatenates each topic's messages and then sorts the
|
||||
// lot by Time, which has second granularity. That sort must not reorder messages that
|
||||
// share a timestamp, or a topic's own messages come back out of publish order. See #1297.
|
||||
//
|
||||
// The messages have to straddle a second boundary: if every timestamp is identical the
|
||||
// concatenation is already sorted and Go's pdqsort leaves it alone, hiding the bug.
|
||||
s := newTestServer(t, newTestConfig(t, databaseURL))
|
||||
|
||||
const perBatch = 10
|
||||
publish := func(batch int) {
|
||||
for _, topic := range []string{"topicA", "topicB"} {
|
||||
for i := 0; i < perBatch; i++ {
|
||||
body := fmt.Sprintf("%s-%02d", topic, batch*perBatch+i)
|
||||
require.Equal(t, 200, request(t, s, "PUT", "/"+topic, body, nil).Code)
|
||||
}
|
||||
}
|
||||
}
|
||||
publish(0)
|
||||
time.Sleep(1100 * time.Millisecond) // Cross a second boundary, so Time is not all-equal
|
||||
publish(1)
|
||||
|
||||
response := request(t, s, "GET", "/topicA,topicB/json?poll=1", "", nil)
|
||||
require.Equal(t, 200, response.Code)
|
||||
messages := toMessages(t, response.Body.String())
|
||||
require.Equal(t, 4*perBatch, len(messages))
|
||||
|
||||
// Each topic's own messages must appear in publish order, whatever the interleaving
|
||||
lastSeen := map[string]int{"topicA": -1, "topicB": -1}
|
||||
for _, m := range messages {
|
||||
topic, seqStr, found := strings.Cut(m.Message, "-")
|
||||
require.True(t, found)
|
||||
seq, err := strconv.Atoi(seqStr)
|
||||
require.Nil(t, err)
|
||||
require.Greater(t, seq, lastSeen[topic], "%s came back out of publish order", m.Message)
|
||||
lastSeen[topic] = seq
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
func TestServer_PollMessageLimit(t *testing.T) {
|
||||
forEachBackend(t, func(t *testing.T, databaseURL string) {
|
||||
// A poll without "since" replays the whole cache, so the response is unbounded unless the
|
||||
// query is capped. The cap keeps the NEWEST messages and advertises the truncation, since
|
||||
// "since" only moves forward and a client cannot fetch what was dropped.
|
||||
c := newTestConfig(t, databaseURL)
|
||||
c.MessagePollLimit = 5
|
||||
s := newTestServer(t, c)
|
||||
|
||||
for i := 1; i <= 8; i++ {
|
||||
require.Equal(t, 200, request(t, s, "PUT", "/mytopic", fmt.Sprintf("message %d", i), nil).Code)
|
||||
}
|
||||
|
||||
response := request(t, s, "GET", "/mytopic/json?poll=1", "", nil)
|
||||
require.Equal(t, 200, response.Code)
|
||||
require.Equal(t, "1", response.Header().Get("X-Messages-Truncated"))
|
||||
messages := toMessages(t, response.Body.String())
|
||||
require.Equal(t, 5, len(messages))
|
||||
require.Equal(t, "message 4", messages[0].Message) // newest 5, oldest first
|
||||
require.Equal(t, "message 8", messages[4].Message)
|
||||
|
||||
// A topic under the limit is served in full, with no truncation header
|
||||
require.Equal(t, 200, request(t, s, "PUT", "/othertopic", "only message", nil).Code)
|
||||
response = request(t, s, "GET", "/othertopic/json?poll=1", "", nil)
|
||||
require.Equal(t, 200, response.Code)
|
||||
require.Empty(t, response.Header().Get("X-Messages-Truncated"))
|
||||
require.Equal(t, 1, len(toMessages(t, response.Body.String())))
|
||||
})
|
||||
}
|
||||
|
||||
func TestServer_PollBandwidthLimit(t *testing.T) {
|
||||
forEachBackend(t, func(t *testing.T, databaseURL string) {
|
||||
// A poll without "since" replays the entire cache, so a topic that is cheap to fill is
|
||||
|
||||
Reference in New Issue
Block a user