From 7826c86f67ff78647b0c69d1717f3e0d34246ec9 Mon Sep 17 00:00:00 2001 From: binwiederhier Date: Thu, 27 Aug 2026 19:42:30 +0000 Subject: [PATCH] Remove message-poll-limit again, hardcode title/tags/etc limits, count polls toward bandwidth limit --- cmd/serve.go | 5 -- docs/config.md | 3 +- docs/publish.md | 1 + docs/releases.md | 3 +- docs/subscribe/api.md | 8 +-- message/cache.go | 73 +++++++++++++++--------- message/cache_postgres.go | 4 -- message/cache_sqlite.go | 4 -- model/model.go | 30 ++++++++++ server/config.go | 18 ++++-- server/errors.go | 2 + server/server.go | 29 +++++++--- server/server.yml | 9 --- server/server_test.go | 116 ++++++++++++++++++++++++++++++++++---- 14 files changed, 224 insertions(+), 81 deletions(-) diff --git a/cmd/serve.go b/cmd/serve.go index d9658a85..07e56867 100644 --- a/cmd/serve.go +++ b/cmd/serve.go @@ -83,7 +83,6 @@ 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"}), @@ -206,7 +205,6 @@ 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") @@ -404,8 +402,6 @@ 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 @@ -528,7 +524,6 @@ 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 968a68d9..58bb7dcd 100644 --- a/docs/config.md +++ b/docs/config.md @@ -2378,14 +2378,13 @@ 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. | | `upstream-base-url` | `NTFY_UPSTREAM_BASE_URL` | *URL* | `https://ntfy.sh` | Forward poll request to an upstream server, this is needed for iOS push notifications for self-hosted servers | | `upstream-access-token` | `NTFY_UPSTREAM_ACCESS_TOKEN` | *string* | `tk_zyYLYj...` | Access token to use for the upstream server; needed only if upstream rate limits are exceeded or upstream server requires auth | | `visitor-attachment-total-size-limit` | `NTFY_VISITOR_ATTACHMENT_TOTAL_SIZE_LIMIT` | *size* | 100M | Rate limiting: Total storage limit used for attachments per visitor, for all attachments combined. Storage is freed after attachments expire. See `attachment-expiry-duration`. | -| `visitor-attachment-daily-bandwidth-limit` | `NTFY_VISITOR_ATTACHMENT_DAILY_BANDWIDTH_LIMIT` | *size* | 500M | Rate limiting: Total daily traffic limit per visitor, covering attachment downloads/uploads and messages replayed from the cache by poll requests. This is to protect your bandwidth costs from exploding. | +| `visitor-attachment-daily-bandwidth-limit` | `NTFY_VISITOR_ATTACHMENT_DAILY_BANDWIDTH_LIMIT` | *size* | 500M | Rate limiting: Total daily traffic limit per visitor, covering attachment downloads/uploads and messages replayed from the cache by poll requests. This is to protect your bandwidth costs from exploding. | | `visitor-email-limit-burst` | `NTFY_VISITOR_EMAIL_LIMIT_BURST` | *number* | 16 | Rate limiting:Initial limit of e-mails per visitor | | `visitor-email-limit-replenish` | `NTFY_VISITOR_EMAIL_LIMIT_REPLENISH` | *duration* | 1h | Rate limiting: Strongly related to `visitor-email-limit-burst`: The rate at which the bucket is refilled | | `visitor-message-daily-limit` | `NTFY_VISITOR_MESSAGE_DAILY_LIMIT` | *number* | - | Rate limiting: Allowed number of messages per day per visitor, reset every day at midnight (UTC). By default, this value is unset. | diff --git a/docs/publish.md b/docs/publish.md index 630d3327..d3bafd22 100644 --- a/docs/publish.md +++ b/docs/publish.md @@ -4931,6 +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. | +| **Title and tag size** | The message title is limited to 1 KB, and all tags combined to 512 bytes. Requests exceeding either are rejected with HTTP 400. | | **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. | diff --git a/docs/releases.md b/docs/releases.md index 79b33d69..ce77f223 100644 --- a/docs/releases.md +++ b/docs/releases.md @@ -2082,7 +2082,8 @@ 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 +* Limit the message title to 1 KB and all tags combined to 512 bytes, rejecting larger requests with HTTP 400 (error codes `40057` and `40058`). Neither field had a size limit before, unlike the message body; on ntfy.sh the 99.9th percentile is 212 bytes for titles and 244 for tags +* Cap a single cache replay at 10 MB of messages per topic. 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, so one request could allocate that much on the server. The newest messages that fit are kept and a truncated response carries an `X-Messages-Truncated: 1` header * `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:** diff --git a/docs/subscribe/api.md b/docs/subscribe/api.md index 3a831028..b1309223 100644 --- a/docs/subscribe/api.md +++ b/docs/subscribe/api.md @@ -283,10 +283,10 @@ Reading cached messages (a `poll=1` request, or any request with `since=`) repla 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. +* **The response is size-capped.** A replay returns only the newest messages that fit in 10 MB + **per topic** (counting body, title, tags and every other publisher-set field), and a capped response carries an `X-Messages-Truncated: 1` header. If you see that + header, older messages were dropped and you did not receive the full cache. In practice this only + affects very large topics; a client polling with `since=` never comes close. * **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. diff --git a/message/cache.go b/message/cache.go index 66d442df..710d5f1e 100644 --- a/message/cache.go +++ b/message/cache.go @@ -20,9 +20,8 @@ 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 + // NoLimit reads a topic's cached messages without a size budget. + NoLimit = 0 ) var errNoRows = errors.New("no rows found") @@ -198,49 +197,49 @@ func (c *Cache) Messages(topic string, since model.SinceMarker, scheduled bool) 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) { +// MessagesCapped returns cached messages for a topic, oldest first, keeping the newest messages +// that fit in maxBytes worth of Message.Size (0 = no budget). The bool reports whether older messages +// were dropped, so the caller can tell the client that what it got is incomplete. +func (c *Cache) MessagesCapped(topic string, since model.SinceMarker, scheduled bool, maxBytes int64) ([]*model.Message, bool, error) { if since.IsNone() { return make([]*model.Message, 0), false, nil } else if since.IsLatest() { messages, err := c.messagesLatest(topic) return messages, false, err } else if since.IsID() { - return c.messagesSinceID(topic, since, scheduled, limit) + return c.messagesSinceID(topic, since, scheduled, maxBytes) } - return c.messagesSinceTime(topic, since, scheduled, limit) + return c.messagesSinceTime(topic, since, scheduled, maxBytes) } -func (c *Cache) messagesSinceTime(topic string, since model.SinceMarker, scheduled bool, limit int) ([]*model.Message, bool, error) { +func (c *Cache) messagesSinceTime(topic string, since model.SinceMarker, scheduled bool, maxBytes int64) ([]*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(), limit+1) + rows, err = rdb.Query(c.queries.selectMessagesSinceTimeScheduled, topic, since.Time().Unix()) } else { - rows, err = rdb.Query(c.queries.selectMessagesSinceTime, topic, since.Time().Unix(), limit+1) + rows, err = rdb.Query(c.queries.selectMessagesSinceTime, topic, since.Time().Unix()) } if err != nil { return nil, false, err } - return readMessagesCapped(rows, limit) + return readMessagesCapped(rows, maxBytes) } -func (c *Cache) messagesSinceID(topic string, since model.SinceMarker, scheduled bool, limit int) ([]*model.Message, bool, error) { +func (c *Cache) messagesSinceID(topic string, since model.SinceMarker, scheduled bool, maxBytes int64) ([]*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(), limit+1) + rows, err = rdb.Query(c.queries.selectMessagesSinceIDScheduled, topic, since.ID()) } else { - rows, err = rdb.Query(c.queries.selectMessagesSinceID, topic, since.ID(), limit+1) + rows, err = rdb.Query(c.queries.selectMessagesSinceID, topic, since.ID()) } if err != nil { return nil, false, err } - return readMessagesCapped(rows, limit) + return readMessagesCapped(rows, maxBytes) } func (c *Cache) messagesLatest(topic string) ([]*model.Message, error) { @@ -473,17 +472,37 @@ 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 +// readMessagesCapped reads a newest-first result set, keeping the newest messages that fit in +// maxBytes worth of Message.Size (0 = no budget), and reverses them into the oldest-first +// order callers expect. It stops scanning once the budget is spent rather than reading everything +// and trimming, so a replay of a huge topic never materializes the whole cache. The bool reports +// whether older messages were left behind. +func readMessagesCapped(rows *sql.Rows, maxBytes int64) ([]*model.Message, bool, error) { + defer rows.Close() + messages := make([]*model.Message, 0) + truncated := false + var total int64 + for rows.Next() { + m, err := readMessage(rows) + if err != nil { + return nil, false, err + } + if maxBytes > 0 { + size := int64(m.Size()) + // Always return at least one message, even if it alone exceeds the budget: an empty + // reply is less useful than an oversized one, and the per-field limits bound how big it gets. + if len(messages) > 0 && total+size > maxBytes { + truncated = true + break + } + total += size + } + messages = append(messages, m) } - truncated := len(messages) > limit - if truncated { - messages = messages[:limit] + if !truncated { + if err := rows.Err(); err != nil { + return nil, false, err + } } slices.Reverse(messages) return messages, truncated, nil diff --git a/message/cache_postgres.go b/message/cache_postgres.go index 3f74dbd8..efecb8de 100644 --- a/message/cache_postgres.go +++ b/message/cache_postgres.go @@ -26,14 +26,12 @@ const ( FROM message WHERE topic = $1 AND time >= $2 AND published = TRUE 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 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 @@ -42,7 +40,6 @@ const ( AND id > COALESCE((SELECT id FROM message WHERE mid = $2), 0) AND published = TRUE 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 @@ -50,7 +47,6 @@ const ( WHERE topic = $1 AND (id > COALESCE((SELECT id FROM message WHERE mid = $2), 0) OR published = FALSE) 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 7c894291..debc495c 100644 --- a/message/cache_sqlite.go +++ b/message/cache_sqlite.go @@ -32,28 +32,24 @@ const ( FROM messages WHERE topic = ? AND time >= ? AND published = 1 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 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 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 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/model/model.go b/model/model.go index b3dae915..e8143eba 100644 --- a/model/model.go +++ b/model/model.go @@ -102,6 +102,24 @@ func (m *Message) ForJSON() *Message { } // Attachment represents a file attachment on a message +// Size returns an approximate byte size of the variable-length, publisher-controlled parts of a +// message. It is used to budget cache replays, so it deliberately counts every field a publisher +// can grow rather than trying to match the exact wire size. +func (m *Message) Size() int { + size := len(m.ID) + len(m.SequenceID) + len(m.Event) + len(m.Topic) + len(m.Title) + + len(m.Message) + len(m.Click) + len(m.Icon) + len(m.ContentType) + len(m.Encoding) + len(m.PollID) + for _, tag := range m.Tags { + size += len(tag) + } + for _, action := range m.Actions { + size += action.Size() + } + if m.Attachment != nil { + size += len(m.Attachment.Name) + len(m.Attachment.Type) + len(m.Attachment.URL) + } + return size +} + type Attachment struct { Name string `json:"name"` Type string `json:"type,omitempty"` @@ -125,6 +143,18 @@ type Action struct { Value string `json:"value,omitempty"` // used in "copy" action } +// Size returns an approximate byte size of an action's variable-length fields. +func (a *Action) Size() int { + size := len(a.ID) + len(a.Action) + len(a.Label) + len(a.URL) + len(a.Method) + len(a.Body) + len(a.Intent) + len(a.Value) + for key, value := range a.Headers { + size += len(key) + len(value) + } + for key, value := range a.Extras { + size += len(key) + len(value) + } + return size +} + // NewAction creates a new action with initialized maps func NewAction() *Action { return &Action{ diff --git a/server/config.go b/server/config.go index 42571ac6..3a6adcc6 100644 --- a/server/config.go +++ b/server/config.go @@ -4,7 +4,6 @@ import ( "crypto/sha256" "encoding/json" "fmt" - "heckel.io/ntfy/v2/message" "io/fs" "net/netip" "reflect" @@ -67,14 +66,23 @@ 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 - DefaultMessagePollLimit = message.NoLimit // Uncapped by default; operators cap it to bound replay memory (limit x MessageSizeLimit) + DefaultMessageSizeLimit = 4096 // Bytes; note that FCM/APNS have a limit of ~4 KB for the entire message DefaultTotalTopicLimit = 15000 DefaultAttachmentTotalSizeLimit = int64(5 * 1024 * 1024 * 1024) // 5 GB DefaultAttachmentFileSizeLimit = int64(15 * 1024 * 1024) // 15 MB DefaultAttachmentExpiryDuration = 3 * time.Hour DefaultAttachmentOrphanGracePeriod = time.Hour // Don't delete orphaned objects younger than this to avoid races with in-flight uploads + // DefaultMessagePollSizeLimit caps what one cache replay returns per topic. It is a backstop + // against a single request materializing an entire topic cache, not a tunable: on ntfy.sh it + // would fire on 2 of ~98k cached topics. See docs/subscribe/api.md#replay-limits. + DefaultMessagePollSizeLimit = 10 * 1024 * 1024 + + // messageTitleSizeLimit and messageTagsSizeLimit cap two publisher-controlled fields that + // otherwise have no limit of their own. Sized off ntfy.sh's own cache: title p999 is 212 bytes + // (16 of ~3M messages exceed 1 KB), tags p999 is 244 (197 exceed 512). + messageTitleSizeLimit = 1024 + messageTagsSizeLimit = 512 ) // Defines all per-visitor limits @@ -176,7 +184,7 @@ type Config struct { MessageDelayMin time.Duration MessageDelayMax time.Duration MessageSizeLimit int - MessagePollLimit int + MessagePollSizeLimit int64 TotalTopicLimit int TotalAttachmentSizeLimit int64 VisitorSubscriptionLimit int @@ -285,7 +293,7 @@ func NewConfig() *Config { TwilioVerifyService: "", TwilioCallFormat: nil, MessageSizeLimit: DefaultMessageSizeLimit, - MessagePollLimit: DefaultMessagePollLimit, + MessagePollSizeLimit: DefaultMessagePollSizeLimit, MessageDelayMin: DefaultMessageDelayMin, MessageDelayMax: DefaultMessageDelayMax, TotalTopicLimit: DefaultTotalTopicLimit, diff --git a/server/errors.go b/server/errors.go index 3e03894f..d798fb5a 100644 --- a/server/errors.go +++ b/server/errors.go @@ -149,6 +149,8 @@ var ( errHTTPBadRequestAnonymousEmailNotAllowed = &errHTTP{40053, http.StatusBadRequest, "invalid request: anonymous email sending is not allowed", "https://ntfy.sh/docs/publish/#e-mail-notifications", nil} errHTTPBadRequestResetLinkInvalid = &errHTTP{40054, http.StatusBadRequest, "invalid request: password reset link invalid or expired", "", nil} errHTTPBadRequestTemplateTooLarge = &errHTTP{40056, http.StatusBadRequest, "invalid request: template too large", "https://ntfy.sh/docs/publish/#message-templating", nil} + errHTTPBadRequestTitleTooLarge = &errHTTP{40057, http.StatusBadRequest, "invalid request: title is too large", "https://ntfy.sh/docs/publish/#limitations", nil} + errHTTPBadRequestTagsTooLarge = &errHTTP{40058, http.StatusBadRequest, "invalid request: tags are too large", "https://ntfy.sh/docs/publish/#limitations", nil} errHTTPNotFound = &errHTTP{40401, http.StatusNotFound, "page not found", "", nil} errHTTPUnauthorized = &errHTTP{40101, http.StatusUnauthorized, "unauthorized", "https://ntfy.sh/docs/publish/#authentication", nil} errHTTPForbidden = &errHTTP{40301, http.StatusForbidden, "forbidden", "https://ntfy.sh/docs/publish/#authentication", nil} diff --git a/server/server.go b/server/server.go index c5b06246..9bd93b06 100644 --- a/server/server.go +++ b/server/server.go @@ -1147,6 +1147,9 @@ func (s *Server) parsePublishParams(r *http.Request, m *model.Message) (cache bo cache = readBoolParam(r, true, "x-cache", "cache") firebase = readBoolParam(r, true, "x-firebase", "firebase") m.Title = readParam(r, "x-title", "title", "t") + if len(m.Title) > messageTitleSizeLimit { + return false, false, "", "", "", false, "", errHTTPBadRequestTitleTooLarge + } m.Click = readParam(r, "x-click", "click") icon := readParam(r, "x-icon", "icon") filename := readParam(r, "x-filename", "filename", "file", "f") @@ -1213,6 +1216,14 @@ func (s *Server) parsePublishParams(r *http.Request, m *model.Message) (cache bo priorityStr = "" // Clear since it's already parsed } m.Tags = readCommaSeparatedParam(r, "x-tags", "tags", "tag", "ta") + // Measured across all tags, not each one: a publisher can add arbitrarily many + tagsSize := 0 + for _, tag := range m.Tags { + tagsSize += len(tag) + } + if tagsSize > messageTagsSizeLimit { + return false, false, "", "", "", false, "", errHTTPBadRequestTagsTooLarge + } delayStr := readParam(r, "x-delay", "delay", "x-at", "at", "x-in", "in") if delayStr != "" { if !cache { @@ -1442,13 +1453,14 @@ func (s *Server) handleSubscribeHTTP(w http.ResponseWriter, r *http.Request, v * if !filters.Pass(msg) { return nil } - m, err := encoder(msg) + encoded, err := encoder(msg) if err != nil { return err } - // Charge before writing, so an exhausted budget fails the first message and surfaces as a - // clean 429 with nothing written. - if meterPollBandwidth && !v.BandwidthAllowed(int64(len(m))) { + // Charge the encoded length, i.e. what actually goes over the wire. Charge before writing, + // so an exhausted budget fails the first message and surfaces as a clean 429 with nothing + // written. + if meterPollBandwidth && !v.BandwidthAllowed(int64(len(encoded))) { return errHTTPTooManyRequestsLimitAttachmentBandwidth } wlock.Lock() @@ -1456,7 +1468,7 @@ func (s *Server) handleSubscribeHTTP(w http.ResponseWriter, r *http.Request, v * if closed { return nil } - if _, err := w.Write([]byte(m)); err != nil { + if _, err := w.Write([]byte(encoded)); err != nil { return err } if fl, ok := w.(http.Flusher); ok { @@ -1741,7 +1753,7 @@ func (s *Server) sendOldMessages(w http.ResponseWriter, topics []*topic, since m messages := make([]*model.Message, 0) truncated := false for _, t := range topics { - topicMessages, topicTruncated, err := s.messageCache.MessagesCapped(t.ID, since, scheduled, s.config.MessagePollLimit) + topicMessages, topicTruncated, err := s.messageCache.MessagesCapped(t.ID, since, scheduled, s.config.MessagePollSizeLimit) if err != nil { return err } @@ -1753,8 +1765,9 @@ func (s *Server) sendOldMessages(w http.ResponseWriter, topics []*topic, since m 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 { + // Must be set before the first message is written, or the header is already on the wire. On the + // WebSocket path the response has been hijacked by then, so this is a no-op there. + if truncated { w.Header().Set("X-Messages-Truncated", "1") } for _, m := range messages { diff --git a/server/server.yml b/server/server.yml index d0c524ec..2e77b97b 100644 --- a/server/server.yml +++ b/server/server.yml @@ -327,15 +327,6 @@ # - 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 3059211d..fa2463f0 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -2852,29 +2852,68 @@ func TestServer_PollOrderAcrossTopics(t *testing.T) { }) } -func TestServer_PollMessageLimit(t *testing.T) { +func TestServer_PublishTitleTooLarge(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. + // Title has no length limit of its own, unlike the body, so it is capped here. Prod p999 + // is 212 bytes and only 16 of ~3M cached messages exceed 1 KB. + s := newTestServer(t, newTestConfig(t, databaseURL)) + + require.Equal(t, 200, request(t, s, "PUT", "/mytopic", "x", map[string]string{ + "Title": strings.Repeat("t", messageTitleSizeLimit), + }).Code) + + response := request(t, s, "PUT", "/mytopic", "x", map[string]string{ + "Title": strings.Repeat("t", messageTitleSizeLimit+1), + }) + require.Equal(t, 400, response.Code) + require.Equal(t, 40057, toHTTPError(t, response.Body.String()).Code) + }) +} + +func TestServer_PublishTagsTooLarge(t *testing.T) { + forEachBackend(t, func(t *testing.T, databaseURL string) { + // Same for tags, measured across all of them: prod p999 is 244 bytes and only 197 of ~3M + // cached messages exceed 512. + s := newTestServer(t, newTestConfig(t, databaseURL)) + + require.Equal(t, 200, request(t, s, "PUT", "/mytopic", "x", map[string]string{ + "Tags": strings.Repeat("g", messageTagsSizeLimit), + }).Code) + + response := request(t, s, "PUT", "/mytopic", "x", map[string]string{ + "Tags": strings.Repeat("g", messageTagsSizeLimit+1), + }) + require.Equal(t, 400, response.Code) + require.Equal(t, 40058, toHTTPError(t, response.Body.String()).Code) + }) +} + +func TestServer_PollSizeLimit(t *testing.T) { + forEachBackend(t, func(t *testing.T, databaseURL string) { + // A poll without "since" replays the entire cache, which is unbounded in size. The cap is + // a byte budget rather than a message count, because message sizes vary ~20x in practice: + // a count cap truncates cheap high-volume topics while barely touching the expensive + // large-message ones it is meant to catch. The newest messages are kept. c := newTestConfig(t, databaseURL) - c.MessagePollLimit = 5 + c.MessagePollSizeLimit = 3500 // Fits three 1000-byte messages, not four 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) + for i := 0; i < 6; i++ { + body := fmt.Sprintf("%04d%s", i, strings.Repeat("x", 996)) // 1000 bytes, ordered prefix + require.Equal(t, 200, request(t, s, "PUT", "/mytopic", body, 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) + require.Equal(t, 3, len(messages)) + require.Equal(t, "0003", messages[0].Message[:4]) // newest three, oldest first + require.Equal(t, "0004", messages[1].Message[:4]) + require.Equal(t, "0005", messages[2].Message[:4]) - // 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) + // A topic under the budget is served whole, with no truncation header + require.Equal(t, 200, request(t, s, "PUT", "/othertopic", "small", 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")) @@ -2882,6 +2921,59 @@ func TestServer_PollMessageLimit(t *testing.T) { }) } +func TestServer_PollSizeLimitCountsTitle(t *testing.T) { + forEachBackend(t, func(t *testing.T, databaseURL string) { + // Title is user-controlled and has no length limit of its own, so it has to count against + // the replay budget too; otherwise a topic of title-heavy messages sails past the cap. + c := newTestConfig(t, databaseURL) + c.MessagePollSizeLimit = 1600 // Fits one 500-byte title + 500-byte body, not two + s := newTestServer(t, c) + + for i := 0; i < 4; i++ { + body := fmt.Sprintf("%04d%s", i, strings.Repeat("b", 496)) // 500 bytes + title := fmt.Sprintf("%04d%s", i, strings.Repeat("t", 496)) // 500 bytes + require.Equal(t, 200, request(t, s, "PUT", "/mytopic", body, map[string]string{"Title": title}).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, 1, len(messages)) // 3 if the title were not counted + require.Equal(t, "0003", messages[0].Message[:4]) + }) +} + +func TestServer_PollSizeLimitCountsEveryField(t *testing.T) { + forEachBackend(t, func(t *testing.T, databaseURL string) { + // Every field a publisher can grow has to count against the replay budget, not just the + // body and title: tags, click, icon and actions are all user-controlled, so anything left + // out is a hole the budget can be walked through. + c := newTestConfig(t, databaseURL) + c.MessagePollSizeLimit = 900 // Two messages fit if only body+title count; one if all fields do + s := newTestServer(t, c) + + tags := make([]string, 5) + for i := range tags { + tags[i] = strings.Repeat("g", 79) // 395 bytes of tags + } + for i := 0; i < 3; i++ { + require.Equal(t, 200, request(t, s, "PUT", "/mytopic", fmt.Sprintf("%04d%s", i, strings.Repeat("b", 196)), map[string]string{ + "Title": strings.Repeat("t", 200), + "Tags": strings.Join(tags, ","), + "Click": "https://example.com/" + strings.Repeat("c", 180), + }).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, 1, len(messages)) // 2 if only body+title were counted + require.Equal(t, "0002", messages[0].Message[:4]) + }) +} + 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