diff --git a/.go-version b/.go-version index 8fe00a57..5db08bf2 100644 --- a/.go-version +++ b/.go-version @@ -1 +1 @@ -1.26.5 +1.27.0 diff --git a/cmd/serve.go b/cmd/serve.go index 07e56867..d9658a85 100644 --- a/cmd/serve.go +++ b/cmd/serve.go @@ -83,6 +83,7 @@ var flagsServe = append( altsrc.NewStringFlag(&cli.StringFlag{Name: "twilio-phone-number", Aliases: []string{"twilio_phone_number"}, EnvVars: []string{"NTFY_TWILIO_PHONE_NUMBER"}, Usage: "Twilio number to use for outgoing calls"}), altsrc.NewStringFlag(&cli.StringFlag{Name: "twilio-verify-service", Aliases: []string{"twilio_verify_service"}, EnvVars: []string{"NTFY_TWILIO_VERIFY_SERVICE"}, Usage: "Twilio Verify service ID, used for phone number verification"}), altsrc.NewStringFlag(&cli.StringFlag{Name: "twilio-call-format", Aliases: []string{"twilio_call_format"}, EnvVars: []string{"NTFY_TWILIO_CALL_FORMAT"}, Usage: "Twilio/TwiML format string for phone calls"}), + altsrc.NewIntFlag(&cli.IntFlag{Name: "message-poll-limit", Aliases: []string{"message_poll_limit"}, EnvVars: []string{"NTFY_MESSAGE_POLL_LIMIT"}, Value: server.DefaultMessagePollLimit, Usage: "max number of cached messages returned per topic when replaying the cache (poll/since requests); default is uncapped"}), altsrc.NewStringFlag(&cli.StringFlag{Name: "message-size-limit", Aliases: []string{"message_size_limit"}, EnvVars: []string{"NTFY_MESSAGE_SIZE_LIMIT"}, Value: util.FormatSize(server.DefaultMessageSizeLimit), Usage: "size limit for the message (see docs for limitations)"}), altsrc.NewStringFlag(&cli.StringFlag{Name: "message-delay-limit", Aliases: []string{"message_delay_limit"}, EnvVars: []string{"NTFY_MESSAGE_DELAY_LIMIT"}, Value: util.FormatDuration(server.DefaultMessageDelayMax), Usage: "max duration a message can be scheduled into the future"}), altsrc.NewIntFlag(&cli.IntFlag{Name: "global-topic-limit", Aliases: []string{"global_topic_limit", "T"}, EnvVars: []string{"NTFY_GLOBAL_TOPIC_LIMIT"}, Value: server.DefaultTotalTopicLimit, Usage: "total number of topics allowed"}), @@ -205,6 +206,7 @@ func execServe(c *cli.Context) error { twilioVerifyService := c.String("twilio-verify-service") twilioCallFormat := c.String("twilio-call-format") messageSizeLimitStr := c.String("message-size-limit") + messagePollLimit := c.Int("message-poll-limit") messageDelayLimitStr := c.String("message-delay-limit") totalTopicLimit := c.Int("global-topic-limit") visitorSubscriptionLimit := c.Int("visitor-subscription-limit") @@ -402,6 +404,8 @@ func execServe(c *cli.Context) error { return fmt.Errorf("if ban-file is set, its directory (%s) must exist", filepath.Dir(banFile)) } else if runtime.GOOS == "windows" && listenUnix != "" { return errors.New("listen-unix is not supported on Windows") + } else if messagePollLimit < 1 { + return errors.New("config option message-poll-limit must be at least 1") } // Backwards compatibility @@ -524,6 +528,7 @@ func execServe(c *cli.Context) error { conf.TwilioVerifyService = twilioVerifyService conf.TwilioCallFormat = twilioCallFormatTemplate conf.MessageSizeLimit = int(messageSizeLimit) + conf.MessagePollLimit = messagePollLimit conf.MessageDelayMax = messageDelayLimit conf.TotalTopicLimit = totalTopicLimit conf.VisitorSubscriptionLimit = visitorSubscriptionLimit diff --git a/docs/config.md b/docs/config.md index 892e99d0..968a68d9 100644 --- a/docs/config.md +++ b/docs/config.md @@ -2378,6 +2378,7 @@ variable before running the `ntfy` command (e.g. `export NTFY_LISTEN_HTTP=:80`). | `twilio-verify-service` | `NTFY_TWILIO_VERIFY_SERVICE` | *string* | - | Twilio Verify service SID, e.g. VA12345beefbeef67890beefbeef122586 | | `keepalive-interval` | `NTFY_KEEPALIVE_INTERVAL` | *duration* | 45s | Interval in which keepalive messages are sent to the client. This is to prevent intermediaries closing the connection for inactivity. Note that the Android app has a hardcoded timeout at 77s, so it should be less than that. | | `manager-interval` | `NTFY_MANAGER_INTERVAL` | *duration* | 1m | Interval in which the manager prunes old messages, deletes topics and prints the stats. | +| `message-poll-limit` | `NTFY_MESSAGE_POLL_LIMIT` | *number* | *(uncapped)* | Max number of cached messages returned per topic when the cache is replayed (a poll, or a `since` request). A poll without a `since` cursor returns a topic's entire cache, which is unbounded in size. When set, the newest messages are kept and the response carries an `X-Messages-Truncated: 1` header. Together with `message-size-limit` this bounds the memory one replay can allocate. | | `message-size-limit` | `NTFY_MESSAGE_SIZE_LIMIT` | *size* | 4K | The size limit for the message body. Please note that this is largely untested, and that FCM/APNS have limits around 4KB. If you increase this size limit, FCM and APNS will NOT work for large messages. | | `message-delay-limit` | `NTFY_MESSAGE_DELAY_LIMIT` | *duration* | 3d | Amount of time a message can be [scheduled](publish.md#scheduled-delivery) into the future when using the `Delay` header | | `global-topic-limit` | `NTFY_GLOBAL_TOPIC_LIMIT` | *number* | 15,000 | Rate limiting: Total number of topics before the server rejects new topics. | diff --git a/docs/publish.md b/docs/publish.md index a95b5458..630d3327 100644 --- a/docs/publish.md +++ b/docs/publish.md @@ -4931,7 +4931,7 @@ but just in case, let's list them all: | **Subscription limit** | By default, the server allows each visitor to keep 30 connections to the server open. | | **Attachment size limit** | By default, the server allows attachments up to 15 MB in size, up to 100 MB in total per visitor and up to 5 GB across all visitors. On ntfy.sh, the attachment size limit is 2 MB, and the per-visitor total is 20 MB. | | **Attachment expiry** | By default, the server deletes attachments after 3 hours and thereby frees up space from the total visitor attachment limit. | -| **Attachment bandwidth** | By default, the server allows 500 MB of GET/PUT/POST traffic for attachments per visitor in a 24 hour period. Traffic exceeding that is rejected. On ntfy.sh, the daily bandwidth limit is 200 MB. | +| **Daily bandwidth** | By default, the server allows 500 MB of traffic per visitor in a 24 hour period, covering attachment GET/PUT/POST traffic and messages replayed from the cache by [poll requests](subscribe/api.md#replay-limits). Traffic exceeding that is rejected. On ntfy.sh, the daily bandwidth limit is 200 MB. | | **Total number of topics** | By default, the server is configured to allow 15,000 topics. The ntfy.sh server has higher limits though. | These limits can be changed on a per-user basis using [tiers](config.md#tiers). If [payments](config.md#payments) are enabled, a user tier can be changed by purchasing diff --git a/docs/releases.md b/docs/releases.md index 1f02ea96..79b33d69 100644 --- a/docs/releases.md +++ b/docs/releases.md @@ -2082,8 +2082,13 @@ and the [ntfy Android app](https://github.com/binwiederhier/ntfy-android/release **Features:** +* Add `message-poll-limit`, an optional cap on the number of cached messages returned by a single replay (uncapped by default, so existing behavior is unchanged). A poll without a `since` cursor returns a topic's entire cache, which was previously unbounded and could reach tens of megabytes on a busy topic; when the cap is set, the newest messages are kept and a truncated response carries an `X-Messages-Truncated: 1` header. Combined with `message-size-limit` this bounds the memory a single replay can allocate * `visitor-attachment-daily-bandwidth-limit` now also covers messages replayed from the message cache by poll requests, not just attachment traffic. A poll without a `since` cursor returns a topic's entire cache, so a topic that is cheap to fill can be re-read for many times its own size; polls beyond the budget are rejected with HTTP 429 (error code 42905) before anything is written. **Note that heavy pollers now consume the same budget as attachment downloads**, so operators serving both may want to raise the limit +**Bug fixes + maintenance:** + +* Fix messages being returned out of publish order when polling or replaying **several topics at once** (`/topic1,topic2/json?poll=1`). `Message.Time` has second granularity, so a multi-topic replay sorts many equal keys; the sort was unstable, which could shuffle a single topic's own messages. Single-topic replays were not affected ([#1297](https://github.com/binwiederhier/ntfy/issues/1297)) + ### ntfy iOS app v1.8.0 (UNRELEASED) **Features:** diff --git a/docs/subscribe/api.md b/docs/subscribe/api.md index a549885b..3a831028 100644 --- a/docs/subscribe/api.md +++ b/docs/subscribe/api.md @@ -245,6 +245,9 @@ combined with `since=` (defaults to `since=all`). curl -s "ntfy.sh/mytopic/json?poll=1" ``` +Note that a poll without `since=` returns a topic's **entire cache**, which on a busy topic can be +large. See [replay limits](#replay-limits) below. + ### Fetch cached messages Messages may be cached for a couple of hours (see [message caching](../config.md#message-cache)) to account for network interruptions of subscribers. If the server has configured message caching, you can read back what you missed by using @@ -275,6 +278,28 @@ parameter (makes most sense with the `poll=1` parameter): curl -s "ntfy.sh/mytopic/json?poll=1&sched=1" ``` +### Replay limits +Reading cached messages (a `poll=1` request, or any request with `since=`) replays messages the server +already stored, so unlike a live subscription its cost grows with the size of the topic's cache. Two +server-side limits apply, both of which a well-behaved client should handle: + +* **The number of messages may be capped.** If the server sets `message-poll-limit`, a replay returns + at most that many messages **per topic**, keeping the newest ones, and the response carries an + `X-Messages-Truncated: 1` header. If you see that header, older messages were dropped: you did not + receive the full cache. It is uncapped by default; ntfy.sh caps it. +* **Replayed bytes count against your daily bandwidth budget**, the same one attachment downloads use + (see [limitations](../publish.md#limitations)). Exceeding it returns `HTTP 429` with ntfy error code + `42905`, and no messages are written. + +Both limits exist because a poll without `since=` re-reads the whole cache every time. If you are +polling repeatedly, **pass `since=`** rather than re-fetching everything. The limits +still apply to a `since=` replay, but it returns only what is new, so in practice you will not come +near either one: + +``` +curl -s "ntfy.sh/mytopic/json?poll=1&since=nFS3knfcQ1xe" +``` + ### Filter messages You can filter which messages are returned based on the well-known message fields `id`, `message`, `title`, `priority` and `tags`. Here's an example that only returns messages of high or urgent priority that contains the both tags @@ -308,6 +333,11 @@ $ curl -s ntfy.sh/mytopic1,mytopic2/json {"id":"Cm02DsxUHb","time":1637182643,"event":"message","topic":"mytopic2","message":"for topic 2"} ``` +When replaying cached messages for several topics at once, they are ordered by their `time` field. +Because `time` has **second granularity**, messages published within the same second share a sort key: +each topic's own messages stay in publish order, but the interleaving *between* topics is not defined. +If you need a total order across topics, sort by `time` and fall back to the order received. + ### Authentication Depending on whether the server is configured to support [access control](../config.md#access-control), some topics may be read/write protected so that only users with the correct credentials can subscribe or publish to them. @@ -427,7 +457,7 @@ and can be passed as **HTTP headers** or **query parameters in the URL**. They a | Parameter | Aliases (case-insensitive) | Description | |-------------|----------------------------|---------------------------------------------------------------------------------| -| `poll` | `X-Poll`, `po` | Return cached messages and close connection | +| `poll` | `X-Poll`, `po` | Return cached messages and close connection (see [replay limits](#replay-limits)) | | `since` | `X-Since`, `si` | Return cached messages since timestamp, duration or message ID | | `scheduled` | `X-Scheduled`, `sched` | Include scheduled/delayed messages in message list | | `id` | `X-ID` | Filter: Only return messages that match this exact message ID | diff --git a/message/cache.go b/message/cache.go index 73aaa076..66d442df 100644 --- a/message/cache.go +++ b/message/cache.go @@ -5,6 +5,7 @@ import ( "encoding/json" "errors" "net/netip" + "slices" "strings" "sync" "time" @@ -18,6 +19,10 @@ import ( const ( tagMessageCache = "message_cache" schemaStore = "message" // Store name in the schema_version table (see db/schema) + + // NoLimit is the message limit for an uncapped read: large enough never to bind in practice, + // small enough that limit+1 stays well inside int32 for the SQL driver. + NoLimit = 1 << 30 ) var errNoRows = errors.New("no rows found") @@ -186,46 +191,56 @@ func (c *Cache) addMessages(ms []*model.Message) error { return nil } -// Messages returns messages for a topic since the given marker, optionally including scheduled messages +// Messages returns all cached messages for a topic, oldest first. Prefer MessagesCapped on +// request paths: an uncapped replay of a busy topic is as large as the topic's entire cache. func (c *Cache) Messages(topic string, since model.SinceMarker, scheduled bool) ([]*model.Message, error) { + messages, _, err := c.MessagesCapped(topic, since, scheduled, NoLimit) + return messages, err +} + +// MessagesCapped returns cached messages for a topic, oldest first, at most limit of them, +// keeping the newest. The bool reports whether older messages were dropped to stay under the +// cap, so the caller can tell the client that what it got is incomplete. +func (c *Cache) MessagesCapped(topic string, since model.SinceMarker, scheduled bool, limit int) ([]*model.Message, bool, error) { if since.IsNone() { - return make([]*model.Message, 0), nil + return make([]*model.Message, 0), false, nil } else if since.IsLatest() { - return c.messagesLatest(topic) + messages, err := c.messagesLatest(topic) + return messages, false, err } else if since.IsID() { - return c.messagesSinceID(topic, since, scheduled) + return c.messagesSinceID(topic, since, scheduled, limit) } - return c.messagesSinceTime(topic, since, scheduled) + return c.messagesSinceTime(topic, since, scheduled, limit) } -func (c *Cache) messagesSinceTime(topic string, since model.SinceMarker, scheduled bool) ([]*model.Message, error) { +func (c *Cache) messagesSinceTime(topic string, since model.SinceMarker, scheduled bool, limit int) ([]*model.Message, bool, error) { var rows *sql.Rows var err error rdb := c.db.ReadOnly() if scheduled { - rows, err = rdb.Query(c.queries.selectMessagesSinceTimeScheduled, topic, since.Time().Unix()) + rows, err = rdb.Query(c.queries.selectMessagesSinceTimeScheduled, topic, since.Time().Unix(), limit+1) } else { - rows, err = rdb.Query(c.queries.selectMessagesSinceTime, topic, since.Time().Unix()) + rows, err = rdb.Query(c.queries.selectMessagesSinceTime, topic, since.Time().Unix(), limit+1) } if err != nil { - return nil, err + return nil, false, err } - return readMessages(rows) + return readMessagesCapped(rows, limit) } -func (c *Cache) messagesSinceID(topic string, since model.SinceMarker, scheduled bool) ([]*model.Message, error) { +func (c *Cache) messagesSinceID(topic string, since model.SinceMarker, scheduled bool, limit int) ([]*model.Message, bool, error) { var rows *sql.Rows var err error rdb := c.db.ReadOnly() if scheduled { - rows, err = rdb.Query(c.queries.selectMessagesSinceIDScheduled, topic, since.ID()) + rows, err = rdb.Query(c.queries.selectMessagesSinceIDScheduled, topic, since.ID(), limit+1) } else { - rows, err = rdb.Query(c.queries.selectMessagesSinceID, topic, since.ID()) + rows, err = rdb.Query(c.queries.selectMessagesSinceID, topic, since.ID(), limit+1) } if err != nil { - return nil, err + return nil, false, err } - return readMessages(rows) + return readMessagesCapped(rows, limit) } func (c *Cache) messagesLatest(topic string) ([]*model.Message, error) { @@ -458,6 +473,22 @@ func (c *Cache) processMessageBatches() { } } +// readMessagesCapped reads a newest-first result set that was queried with LIMIT limit+1. The +// extra row is how overflow is detected without a second COUNT query: if it is there, older +// messages exist beyond the cap. Rows are reversed into the oldest-first order callers expect. +func readMessagesCapped(rows *sql.Rows, limit int) ([]*model.Message, bool, error) { + messages, err := readMessages(rows) + if err != nil { + return nil, false, err + } + truncated := len(messages) > limit + if truncated { + messages = messages[:limit] + } + slices.Reverse(messages) + return messages, truncated, nil +} + func readMessages(rows *sql.Rows) ([]*model.Message, error) { defer rows.Close() messages := make([]*model.Message, 0) diff --git a/message/cache_postgres.go b/message/cache_postgres.go index e588f8c4..3f74dbd8 100644 --- a/message/cache_postgres.go +++ b/message/cache_postgres.go @@ -25,13 +25,15 @@ const ( SELECT mid, sequence_id, time, event, expires, topic, message, title, priority, tags, click, icon, actions, attachment_name, attachment_type, attachment_size, attachment_expires, attachment_url, sender, user_id, content_type, encoding FROM message WHERE topic = $1 AND time >= $2 AND published = TRUE - ORDER BY time, id + ORDER BY time DESC, id DESC + LIMIT $3 ` postgresSelectMessagesSinceTimeIncludeScheduledQuery = ` SELECT mid, sequence_id, time, event, expires, topic, message, title, priority, tags, click, icon, actions, attachment_name, attachment_type, attachment_size, attachment_expires, attachment_url, sender, user_id, content_type, encoding FROM message WHERE topic = $1 AND time >= $2 - ORDER BY time, id + ORDER BY time DESC, id DESC + LIMIT $3 ` postgresSelectMessagesSinceIDQuery = ` SELECT mid, sequence_id, time, event, expires, topic, message, title, priority, tags, click, icon, actions, attachment_name, attachment_type, attachment_size, attachment_expires, attachment_url, sender, user_id, content_type, encoding @@ -39,14 +41,16 @@ const ( WHERE topic = $1 AND id > COALESCE((SELECT id FROM message WHERE mid = $2), 0) AND published = TRUE - ORDER BY time, id + ORDER BY time DESC, id DESC + LIMIT $3 ` postgresSelectMessagesSinceIDIncludeScheduledQuery = ` SELECT mid, sequence_id, time, event, expires, topic, message, title, priority, tags, click, icon, actions, attachment_name, attachment_type, attachment_size, attachment_expires, attachment_url, sender, user_id, content_type, encoding FROM message WHERE topic = $1 AND (id > COALESCE((SELECT id FROM message WHERE mid = $2), 0) OR published = FALSE) - ORDER BY time, id + ORDER BY time DESC, id DESC + LIMIT $3 ` postgresSelectMessagesLatestQuery = ` SELECT mid, sequence_id, time, event, expires, topic, message, title, priority, tags, click, icon, actions, attachment_name, attachment_type, attachment_size, attachment_expires, attachment_url, sender, user_id, content_type, encoding diff --git a/message/cache_sqlite.go b/message/cache_sqlite.go index cc01b537..7c894291 100644 --- a/message/cache_sqlite.go +++ b/message/cache_sqlite.go @@ -31,25 +31,29 @@ const ( SELECT mid, sequence_id, time, event, expires, topic, message, title, priority, tags, click, icon, actions, attachment_name, attachment_type, attachment_size, attachment_expires, attachment_url, sender, user, content_type, encoding FROM messages WHERE topic = ? AND time >= ? AND published = 1 - ORDER BY time, id + ORDER BY time DESC, id DESC + LIMIT ? ` sqliteSelectMessagesSinceTimeIncludeScheduledQuery = ` SELECT mid, sequence_id, time, event, expires, topic, message, title, priority, tags, click, icon, actions, attachment_name, attachment_type, attachment_size, attachment_expires, attachment_url, sender, user, content_type, encoding FROM messages WHERE topic = ? AND time >= ? - ORDER BY time, id + ORDER BY time DESC, id DESC + LIMIT ? ` sqliteSelectMessagesSinceIDQuery = ` SELECT mid, sequence_id, time, event, expires, topic, message, title, priority, tags, click, icon, actions, attachment_name, attachment_type, attachment_size, attachment_expires, attachment_url, sender, user, content_type, encoding FROM messages WHERE topic = ? AND id > COALESCE((SELECT id FROM messages WHERE mid = ?), 0) AND published = 1 - ORDER BY time, id + ORDER BY time DESC, id DESC + LIMIT ? ` sqliteSelectMessagesSinceIDIncludeScheduledQuery = ` SELECT mid, sequence_id, time, event, expires, topic, message, title, priority, tags, click, icon, actions, attachment_name, attachment_type, attachment_size, attachment_expires, attachment_url, sender, user, content_type, encoding FROM messages WHERE topic = ? AND (id > COALESCE((SELECT id FROM messages WHERE mid = ?), 0) OR published = 0) - ORDER BY time, id + ORDER BY time DESC, id DESC + LIMIT ? ` sqliteSelectMessagesLatestQuery = ` SELECT mid, sequence_id, time, event, expires, topic, message, title, priority, tags, click, icon, actions, attachment_name, attachment_type, attachment_size, attachment_expires, attachment_url, sender, user, content_type, encoding diff --git a/server/config.go b/server/config.go index b2c169a4..42571ac6 100644 --- a/server/config.go +++ b/server/config.go @@ -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, diff --git a/server/server.go b/server/server.go index 35e69296..c5b06246 100644 --- a/server/server.go +++ b/server/server.go @@ -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 diff --git a/server/server.yml b/server/server.yml index 2e77b97b..d0c524ec 100644 --- a/server/server.yml +++ b/server/server.yml @@ -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. diff --git a/server/server_test.go b/server/server_test.go index 06b68629..3059211d 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -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 diff --git a/template/gotext/GENERATED_FROM b/template/gotext/GENERATED_FROM index 9c89591a..efdc29c4 100644 --- a/template/gotext/GENERATED_FROM +++ b/template/gotext/GENERATED_FROM @@ -1 +1 @@ -go1.26.5 +go1.27.0 diff --git a/template/gotext/template.go b/template/gotext/template.go index f0e70efa..31b0d409 100644 --- a/template/gotext/template.go +++ b/template/gotext/template.go @@ -167,11 +167,13 @@ func (t *Template) Delims(left, right string) *Template { } // Funcs adds the elements of the argument map to the template's function map. -// It must be called before the template is parsed. +// Any function used in the template must be added before the template is +// parsed. Funcs may be called more than once, including after parsing (for +// example, after [Template.Clone]), to replace a function of the same name; +// the replacement is used when the template is executed. // It panics if a value in the map is not a function with appropriate return // type or if the name cannot be used syntactically as a function in a template. -// It is legal to overwrite elements of the map. The return value is the template, -// so calls can be chained. +// The return value is the template, so calls can be chained. func (t *Template) Funcs(funcMap FuncMap) *Template { t.init() t.muFuncs.Lock()