Compare commits

...
17 changed files with 444 additions and 92 deletions
+1 -1
View File
@@ -1 +1 @@
1.26.5
1.27.0
+1 -1
View File
@@ -2384,7 +2384,7 @@ variable before running the `ntfy` command (e.g. `export NTFY_LISTEN_HTTP=:80`).
| `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. |
+38 -38
View File
@@ -34,37 +34,37 @@ as a service starting at boot time.
=== "x86_64/amd64"
```bash
wget https://github.com/binwiederhier/ntfy/releases/download/v2.27.0/ntfy_2.27.0_linux_amd64.tar.gz
tar zxvf ntfy_2.27.0_linux_amd64.tar.gz
sudo cp -a ntfy_2.27.0_linux_amd64/ntfy /usr/local/bin/ntfy
sudo mkdir /etc/ntfy && sudo cp ntfy_2.27.0_linux_amd64/{client,server}/*.yml /etc/ntfy
wget https://github.com/binwiederhier/ntfy/releases/download/v2.28.0/ntfy_2.28.0_linux_amd64.tar.gz
tar zxvf ntfy_2.28.0_linux_amd64.tar.gz
sudo cp -a ntfy_2.28.0_linux_amd64/ntfy /usr/local/bin/ntfy
sudo mkdir /etc/ntfy && sudo cp ntfy_2.28.0_linux_amd64/{client,server}/*.yml /etc/ntfy
sudo ntfy serve
```
=== "armv6"
```bash
wget https://github.com/binwiederhier/ntfy/releases/download/v2.27.0/ntfy_2.27.0_linux_armv6.tar.gz
tar zxvf ntfy_2.27.0_linux_armv6.tar.gz
sudo cp -a ntfy_2.27.0_linux_armv6/ntfy /usr/bin/ntfy
sudo mkdir /etc/ntfy && sudo cp ntfy_2.27.0_linux_armv6/{client,server}/*.yml /etc/ntfy
wget https://github.com/binwiederhier/ntfy/releases/download/v2.28.0/ntfy_2.28.0_linux_armv6.tar.gz
tar zxvf ntfy_2.28.0_linux_armv6.tar.gz
sudo cp -a ntfy_2.28.0_linux_armv6/ntfy /usr/bin/ntfy
sudo mkdir /etc/ntfy && sudo cp ntfy_2.28.0_linux_armv6/{client,server}/*.yml /etc/ntfy
sudo ntfy serve
```
=== "armv7/armhf"
```bash
wget https://github.com/binwiederhier/ntfy/releases/download/v2.27.0/ntfy_2.27.0_linux_armv7.tar.gz
tar zxvf ntfy_2.27.0_linux_armv7.tar.gz
sudo cp -a ntfy_2.27.0_linux_armv7/ntfy /usr/bin/ntfy
sudo mkdir /etc/ntfy && sudo cp ntfy_2.27.0_linux_armv7/{client,server}/*.yml /etc/ntfy
wget https://github.com/binwiederhier/ntfy/releases/download/v2.28.0/ntfy_2.28.0_linux_armv7.tar.gz
tar zxvf ntfy_2.28.0_linux_armv7.tar.gz
sudo cp -a ntfy_2.28.0_linux_armv7/ntfy /usr/bin/ntfy
sudo mkdir /etc/ntfy && sudo cp ntfy_2.28.0_linux_armv7/{client,server}/*.yml /etc/ntfy
sudo ntfy serve
```
=== "arm64"
```bash
wget https://github.com/binwiederhier/ntfy/releases/download/v2.27.0/ntfy_2.27.0_linux_arm64.tar.gz
tar zxvf ntfy_2.27.0_linux_arm64.tar.gz
sudo cp -a ntfy_2.27.0_linux_arm64/ntfy /usr/bin/ntfy
sudo mkdir /etc/ntfy && sudo cp ntfy_2.27.0_linux_arm64/{client,server}/*.yml /etc/ntfy
wget https://github.com/binwiederhier/ntfy/releases/download/v2.28.0/ntfy_2.28.0_linux_arm64.tar.gz
tar zxvf ntfy_2.28.0_linux_arm64.tar.gz
sudo cp -a ntfy_2.28.0_linux_arm64/ntfy /usr/bin/ntfy
sudo mkdir /etc/ntfy && sudo cp ntfy_2.28.0_linux_arm64/{client,server}/*.yml /etc/ntfy
sudo ntfy serve
```
@@ -84,25 +84,25 @@ Install the ntfy server unit file (which contains parameters to start the servic
=== "x86_64/amd64"
```bash
sudo mv ntfy_2.27.0_linux_amd64/server/ntfy.service /etc/systemd/system/
sudo mv ntfy_2.28.0_linux_amd64/server/ntfy.service /etc/systemd/system/
sudo chmod 644 /etc/systemd/system/ntfy.service
```
=== "armv6"
```bash
sudo mv ntfy_2.27.0_linux_armv6/server/ntfy.service /etc/systemd/system/
sudo mv ntfy_2.28.0_linux_armv6/server/ntfy.service /etc/systemd/system/
sudo chmod 644 /etc/systemd/system/ntfy.service
```
=== "armv7/armhf"
```bash
sudo mv ntfy_2.27.0_linux_armv7/server/ntfy.service /etc/systemd/system/
sudo mv ntfy_2.28.0_linux_armv7/server/ntfy.service /etc/systemd/system/
sudo chmod 644 /etc/systemd/system/ntfy.service
```
=== "arm64"
```bash
sudo mv ntfy_2.27.0_linux_arm64/server/ntfy.service /etc/systemd/system/
sudo mv ntfy_2.28.0_linux_arm64/server/ntfy.service /etc/systemd/system/
sudo chmod 644 /etc/systemd/system/ntfy.service
```
@@ -118,25 +118,25 @@ Install the ntfy server service script:
=== "x86_64/amd64"
```bash
sudo mv ntfy_2.27.0_linux_amd64/server/ntfy.openrc /etc/init.d/ntfy
sudo mv ntfy_2.28.0_linux_amd64/server/ntfy.openrc /etc/init.d/ntfy
sudo chmod 755 /etc/init.d/ntfy
```
=== "armv6"
```bash
sudo mv ntfy_2.27.0_linux_armv6/server/ntfy.openrc /etc/init.d/ntfy
sudo mv ntfy_2.28.0_linux_armv6/server/ntfy.openrc /etc/init.d/ntfy
sudo chmod 755 /etc/init.d/ntfy
```
=== "armv7/armhf"
```bash
sudo mv ntfy_2.27.0_linux_armv7/server/ntfy.openrc /etc/init.d/ntfy
sudo mv ntfy_2.28.0_linux_armv7/server/ntfy.openrc /etc/init.d/ntfy
sudo chmod 755 /etc/init.d/ntfy
```
=== "arm64"
```bash
sudo mv ntfy_2.27.0_linux_arm64/server/ntfy.openrc /etc/init.d/ntfy
sudo mv ntfy_2.28.0_linux_arm64/server/ntfy.openrc /etc/init.d/ntfy
sudo chmod 755 /etc/init.d/ntfy
```
@@ -204,7 +204,7 @@ Manually installing the .deb file:
=== "x86_64/amd64"
```bash
wget https://github.com/binwiederhier/ntfy/releases/download/v2.27.0/ntfy_2.27.0_linux_amd64.deb
wget https://github.com/binwiederhier/ntfy/releases/download/v2.28.0/ntfy_2.28.0_linux_amd64.deb
sudo dpkg -i ntfy_*.deb
sudo systemctl enable ntfy
sudo systemctl start ntfy
@@ -212,7 +212,7 @@ Manually installing the .deb file:
=== "armv6"
```bash
wget https://github.com/binwiederhier/ntfy/releases/download/v2.27.0/ntfy_2.27.0_linux_armv6.deb
wget https://github.com/binwiederhier/ntfy/releases/download/v2.28.0/ntfy_2.28.0_linux_armv6.deb
sudo dpkg -i ntfy_*.deb
sudo systemctl enable ntfy
sudo systemctl start ntfy
@@ -220,7 +220,7 @@ Manually installing the .deb file:
=== "armv7/armhf"
```bash
wget https://github.com/binwiederhier/ntfy/releases/download/v2.27.0/ntfy_2.27.0_linux_armv7.deb
wget https://github.com/binwiederhier/ntfy/releases/download/v2.28.0/ntfy_2.28.0_linux_armv7.deb
sudo dpkg -i ntfy_*.deb
sudo systemctl enable ntfy
sudo systemctl start ntfy
@@ -228,7 +228,7 @@ Manually installing the .deb file:
=== "arm64"
```bash
wget https://github.com/binwiederhier/ntfy/releases/download/v2.27.0/ntfy_2.27.0_linux_arm64.deb
wget https://github.com/binwiederhier/ntfy/releases/download/v2.28.0/ntfy_2.28.0_linux_arm64.deb
sudo dpkg -i ntfy_*.deb
sudo systemctl enable ntfy
sudo systemctl start ntfy
@@ -238,28 +238,28 @@ Manually installing the .deb file:
=== "x86_64/amd64"
```bash
sudo rpm -ivh https://github.com/binwiederhier/ntfy/releases/download/v2.27.0/ntfy_2.27.0_linux_amd64.rpm
sudo rpm -ivh https://github.com/binwiederhier/ntfy/releases/download/v2.28.0/ntfy_2.28.0_linux_amd64.rpm
sudo systemctl enable ntfy
sudo systemctl start ntfy
```
=== "armv6"
```bash
sudo rpm -ivh https://github.com/binwiederhier/ntfy/releases/download/v2.27.0/ntfy_2.27.0_linux_armv6.rpm
sudo rpm -ivh https://github.com/binwiederhier/ntfy/releases/download/v2.28.0/ntfy_2.28.0_linux_armv6.rpm
sudo systemctl enable ntfy
sudo systemctl start ntfy
```
=== "armv7/armhf"
```bash
sudo rpm -ivh https://github.com/binwiederhier/ntfy/releases/download/v2.27.0/ntfy_2.27.0_linux_armv7.rpm
sudo rpm -ivh https://github.com/binwiederhier/ntfy/releases/download/v2.28.0/ntfy_2.28.0_linux_armv7.rpm
sudo systemctl enable ntfy
sudo systemctl start ntfy
```
=== "arm64"
```bash
sudo rpm -ivh https://github.com/binwiederhier/ntfy/releases/download/v2.27.0/ntfy_2.27.0_linux_arm64.rpm
sudo rpm -ivh https://github.com/binwiederhier/ntfy/releases/download/v2.28.0/ntfy_2.28.0_linux_arm64.rpm
sudo systemctl enable ntfy
sudo systemctl start ntfy
```
@@ -301,18 +301,18 @@ pkg install go-ntfy
## macOS
The [ntfy CLI](subscribe/cli.md) (`ntfy publish` and `ntfy subscribe` only) is supported on macOS as well.
To install, please [download the tarball](https://github.com/binwiederhier/ntfy/releases/download/v2.27.0/ntfy_2.27.0_darwin_all.tar.gz),
To install, please [download the tarball](https://github.com/binwiederhier/ntfy/releases/download/v2.28.0/ntfy_2.28.0_darwin_all.tar.gz),
extract it and place it somewhere in your `PATH` (e.g. `/usr/local/bin/ntfy`).
If run as `root`, ntfy will look for its config at `/etc/ntfy/client.yml`. For all other users, it'll look for it at
`~/Library/Application Support/ntfy/client.yml` (sample included in the tarball).
```bash
curl -L https://github.com/binwiederhier/ntfy/releases/download/v2.27.0/ntfy_2.27.0_darwin_all.tar.gz > ntfy_2.27.0_darwin_all.tar.gz
tar zxvf ntfy_2.27.0_darwin_all.tar.gz
sudo cp -a ntfy_2.27.0_darwin_all/ntfy /usr/local/bin/ntfy
curl -L https://github.com/binwiederhier/ntfy/releases/download/v2.28.0/ntfy_2.28.0_darwin_all.tar.gz > ntfy_2.28.0_darwin_all.tar.gz
tar zxvf ntfy_2.28.0_darwin_all.tar.gz
sudo cp -a ntfy_2.28.0_darwin_all/ntfy /usr/local/bin/ntfy
mkdir ~/Library/Application\ Support/ntfy
cp ntfy_2.27.0_darwin_all/client/client.yml ~/Library/Application\ Support/ntfy/client.yml
cp ntfy_2.28.0_darwin_all/client/client.yml ~/Library/Application\ Support/ntfy/client.yml
ntfy --help
```
@@ -333,7 +333,7 @@ brew install ntfy
The ntfy server and CLI are fully supported on Windows. You can run the ntfy server directly or as a Windows service.
To install, you can either
* [Download the latest ZIP](https://github.com/binwiederhier/ntfy/releases/download/v2.27.0/ntfy_2.27.0_windows_amd64.zip),
* [Download the latest ZIP](https://github.com/binwiederhier/ntfy/releases/download/v2.28.0/ntfy_2.28.0_windows_amd64.zip),
extract it and place the `ntfy.exe` binary somewhere in your `%Path%`.
* Or install ntfy from the [Scoop](https://scoop.sh) main repository via `scoop install ntfy`
+3 -2
View File
@@ -666,7 +666,7 @@ them with a comma, e.g. `tag1,tag2,tag3`.
_Supported on:_ :material-android: :material-firefox:
You can format messages using [Markdown](https://www.markdownguide.org/basic-syntax/) 🤩. That means you can use
**bold text**, *italicized text*, links, images, and more. Supported Markdown features (web app only for now):
**bold text**, *italicized text*, links, images, and more. Supported Markdown features:
- [Emphasis](https://www.markdownguide.org/basic-syntax/#emphasis) such as **bold** (`**bold**`), *italics* (`*italics*`)
- [Links](https://www.markdownguide.org/basic-syntax/#links) (`[some tool](https://ntfy.sh)`)
@@ -4931,7 +4931,8 @@ 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. |
| **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. |
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
+16 -7
View File
@@ -6,12 +6,27 @@ and the [ntfy Android app](https://github.com/binwiederhier/ntfy-android/release
| Component | Version | Release date |
|------------------|---------|---------------|
| ntfy server | v2.27.0 | Aug 4, 2026 |
| ntfy server | v2.28.0 | Aug 27, 2026 |
| ntfy Android app | v1.25.2 | July 23, 2026 |
| ntfy iOS app | v1.7.0 | May 30, 2026 |
Please check out the release notes for [upcoming releases](#not-released-yet) below.
### ntfy server v2.28.0
Released August 27, 2026
This is a hardening release. A single topic on ntfy.sh was polled continuously with `poll=1` and no
`since` cursor, which replays a topic's entire cache on every request. The changes below bound what one
replay can cost, close two fields that had no size limit at all, and fix an ordering bug found while
digging into it.
**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))
* 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
### ntfy server v2.27.0
Released August 4, 2026
@@ -2078,12 +2093,6 @@ and the [ntfy Android app](https://github.com/binwiederhier/ntfy-android/release
## Not released yet
### ntfy server v2.28.0 (UNRELEASED)
**Features:**
* `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
### ntfy iOS app v1.8.0 (UNRELEASED)
**Features:**
+31 -1
View File
@@ -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 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.
Both limits exist because a poll without `since=` re-reads the whole cache every time. If you are
polling repeatedly, **pass `since=<last message ID>`** 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 |
+67 -16
View File
@@ -5,6 +5,7 @@ import (
"encoding/json"
"errors"
"net/netip"
"slices"
"strings"
"sync"
"time"
@@ -18,6 +19,9 @@ import (
const (
tagMessageCache = "message_cache"
schemaStore = "message" // Store name in the schema_version table (see db/schema)
// NoLimit reads a topic's cached messages without a size budget.
NoLimit = 0
)
var errNoRows = errors.New("no rows found")
@@ -186,19 +190,29 @@ 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) {
if since.IsNone() {
return make([]*model.Message, 0), nil
} else if since.IsLatest() {
return c.messagesLatest(topic)
} else if since.IsID() {
return c.messagesSinceID(topic, since, scheduled)
}
return c.messagesSinceTime(topic, since, scheduled)
messages, _, err := c.MessagesCapped(topic, since, scheduled, NoLimit)
return messages, err
}
func (c *Cache) messagesSinceTime(topic string, since model.SinceMarker, scheduled bool) ([]*model.Message, 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, maxBytes)
}
return c.messagesSinceTime(topic, since, scheduled, maxBytes)
}
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()
@@ -208,12 +222,12 @@ func (c *Cache) messagesSinceTime(topic string, since model.SinceMarker, schedul
rows, err = rdb.Query(c.queries.selectMessagesSinceTime, topic, since.Time().Unix())
}
if err != nil {
return nil, err
return nil, false, err
}
return readMessages(rows)
return readMessagesCapped(rows, maxBytes)
}
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, maxBytes int64) ([]*model.Message, bool, error) {
var rows *sql.Rows
var err error
rdb := c.db.ReadOnly()
@@ -223,9 +237,9 @@ func (c *Cache) messagesSinceID(topic string, since model.SinceMarker, scheduled
rows, err = rdb.Query(c.queries.selectMessagesSinceID, topic, since.ID())
}
if err != nil {
return nil, err
return nil, false, err
}
return readMessages(rows)
return readMessagesCapped(rows, maxBytes)
}
func (c *Cache) messagesLatest(topic string) ([]*model.Message, error) {
@@ -286,7 +300,8 @@ func (c *Cache) MarkPublished(m *model.Message) error {
return err
}
// MessagesCount returns the total number of messages in the cache
// MessagesCount returns the total number of messages in the cache. On Postgres, this is the
// planner's estimate once the table has been analyzed, not an exact count.
func (c *Cache) MessagesCount() (int, error) {
rows, err := c.db.ReadOnly().Query(c.queries.selectMessagesCount)
if err != nil {
@@ -458,6 +473,42 @@ func (c *Cache) processMessageBatches() {
}
}
// 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)
}
if !truncated {
if err := rows.Err(); err != nil {
return nil, false, err
}
}
slices.Reverse(messages)
return messages, truncated, nil
}
func readMessages(rows *sql.Rows) ([]*model.Message, error) {
defer rows.Close()
messages := make([]*model.Message, 0)
+11 -6
View File
@@ -25,13 +25,13 @@ 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
`
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
`
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 +39,14 @@ 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
`
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
`
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
@@ -62,8 +62,13 @@ const (
ORDER BY time, id
`
postgresUpdateMessagePublishedQuery = `UPDATE message SET published = TRUE WHERE mid = $1`
postgresSelectMessagesCountQuery = `SELECT COUNT(*) FROM message`
postgresSelectTopicsQuery = `SELECT topic FROM message GROUP BY topic`
// Planner estimate, since a COUNT(*) scans the whole table; reltuples is -1 if never analyzed
postgresSelectMessagesCountQuery = `
SELECT CASE WHEN reltuples < 0 THEN (SELECT COUNT(*) FROM message) ELSE reltuples::BIGINT END
FROM pg_class
WHERE oid = 'message'::regclass
`
postgresSelectTopicsQuery = `SELECT topic FROM message GROUP BY topic`
postgresDeleteExpiredMessagesQuery = `DELETE FROM message WHERE mid IN (SELECT mid FROM message WHERE expires <= $1 AND published = TRUE LIMIT $2)`
postgresMarkExpiredAttachmentsDeletedQuery = `UPDATE message SET attachment_deleted = TRUE WHERE mid IN (SELECT mid FROM message WHERE attachment_expires > 0 AND attachment_expires <= $1 AND attachment_deleted = FALSE LIMIT $2)`
+25
View File
@@ -74,3 +74,28 @@ func TestPostgresStore_Migration_From14(t *testing.T) {
require.Nil(t, err)
require.Equal(t, dbtest.PostgresSchema(t, freshDB), dbtest.PostgresSchema(t, testDB))
}
func TestPostgresStore_MessagesCount_UsesPlannerEstimate(t *testing.T) {
// The manager calls MessagesCount every minute for a metric; a COUNT(*) scans the whole
// table on every call, so once the table has been analyzed, the planner's estimate is used
testDB := dbtest.CreateTestPostgres(t)
store, err := message.NewPostgresStore(testDB, 0, 0)
require.Nil(t, err)
for i := 0; i < 10; i++ {
require.Nil(t, store.AddMessage(model.NewDefaultMessage("mytopic", "some message")))
}
// Never analyzed: falls back to an exact count
count, err := store.MessagesCount()
require.Nil(t, err)
require.Equal(t, 10, count)
// Analyzed, then rows deleted: the estimate lags until the next (auto)analyze
_, err = testDB.Exec(`ANALYZE message`)
require.Nil(t, err)
_, err = testDB.Exec(`DELETE FROM message WHERE id IN (SELECT id FROM message LIMIT 4)`)
require.Nil(t, err)
count, err = store.MessagesCount()
require.Nil(t, err)
require.Equal(t, 10, count)
}
+4 -4
View File
@@ -31,25 +31,25 @@ 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
`
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
`
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
`
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
`
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
+30
View File
@@ -101,6 +101,24 @@ func (m *Message) ForJSON() *Message {
return m
}
// 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
}
// Attachment represents a file attachment on a message
type Attachment struct {
Name string `json:"name"`
@@ -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{
+12
View File
@@ -73,6 +73,16 @@ const (
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
@@ -174,6 +184,7 @@ type Config struct {
MessageDelayMin time.Duration
MessageDelayMax time.Duration
MessageSizeLimit int
MessagePollSizeLimit int64
TotalTopicLimit int
TotalAttachmentSizeLimit int64
VisitorSubscriptionLimit int
@@ -282,6 +293,7 @@ func NewConfig() *Config {
TwilioVerifyService: "",
TwilioCallFormat: nil,
MessageSizeLimit: DefaultMessageSizeLimit,
MessagePollSizeLimit: DefaultMessagePollSizeLimit,
MessageDelayMin: DefaultMessageDelayMin,
MessageDelayMax: DefaultMessageDelayMax,
TotalTopicLimit: DefaultTotalTopicLimit,
+2
View File
@@ -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}
+33 -12
View File
@@ -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 {
@@ -1474,7 +1486,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 +1502,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 +1637,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 +1651,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 +1746,30 @@ 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.MessagePollSizeLimit)
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. 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 {
if err := sub(v, m); err != nil {
return err
+164
View File
@@ -16,6 +16,7 @@ import (
"os"
"path/filepath"
"runtime/debug"
"strconv"
"strings"
"sync"
"sync/atomic"
@@ -2810,6 +2811,169 @@ 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_PublishTitleTooLarge(t *testing.T) {
forEachBackend(t, func(t *testing.T, databaseURL string) {
// 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.MessagePollSizeLimit = 3500 // Fits three 1000-byte messages, not four
s := newTestServer(t, c)
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, 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 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"))
require.Equal(t, 1, len(toMessages(t, response.Body.String())))
})
}
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
+1 -1
View File
@@ -1 +1 @@
go1.26.5
go1.27.0
+5 -3
View File
@@ -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()