mirror of
https://github.com/multipleof4/ntfy.git
synced 2026-10-09 05:15:22 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e2f4fc7469 | ||
|
|
4f52663dda | ||
|
|
10cb6506f8 | ||
|
|
7826c86f67 | ||
|
|
cae7816e5d |
+1
-1
@@ -1 +1 @@
|
||||
1.26.5
|
||||
1.27.0
|
||||
|
||||
+1
-1
@@ -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
@@ -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
@@ -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
@@ -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
@@ -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
@@ -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)
|
||||
|
||||
@@ -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)`
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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{
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
@@ -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
|
||||
|
||||
@@ -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 @@
|
||||
go1.26.5
|
||||
go1.27.0
|
||||
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user