Compare commits
16 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| fc0d6bf98b | |||
| a44e2b4235 | |||
| 2098343621 | |||
| 5684a02e79 | |||
| 50c3b9db64 | |||
| 39ae43c80e | |||
| c9964778cb | |||
| a553d050ca | |||
| 9c2f3eaf1b | |||
| c63f70ad0c | |||
| bbdcbe213f | |||
| 8c4a5f396b | |||
| b43ba89778 | |||
| f1d360420c | |||
| 697e907c86 | |||
| 894dbbe31d |
27
CHANGELOG.md
27
CHANGELOG.md
@@ -7,6 +7,29 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
|||||||
|
|
||||||
## [Unreleased]
|
## [Unreleased]
|
||||||
|
|
||||||
|
## [0.10.2] - 2026-10-05
|
||||||
|
|
||||||
|
### Fixed
|
||||||
|
|
||||||
|
- A publisher's correction to an item already in iPX, a retitled episode, mended text, new artwork or a fixed length, arrives on the next read. Since 0.9.1 an item kept what it first said: Ain't It Cool News's items each showed the text of the one before, long after the feed was fixed.
|
||||||
|
- A feed's page in the Directory says when the feed is failing, and why, and no longer says nobody here subscribes to a feed someone does.
|
||||||
|
|
||||||
|
## [0.10.1] - 2026-10-05
|
||||||
|
|
||||||
|
### Fixed
|
||||||
|
|
||||||
|
- A feed whose address answers with nothing, as a lapsed domain does behind a DNS filter's block page, shows as failing, with what to do about it, instead of as a show that has not posted yet.
|
||||||
|
- The search box finds episodes in Currently Listening. It did nothing there.
|
||||||
|
- `ipx add`, `ipx rm` and `ipx import`, run while the daemon runs, reach it at once. It kept its old list of feeds and wrote it back at its next change, undoing them.
|
||||||
|
- On SQLite, several requests are answered at once instead of one at a time. With 25 people browsing, a page of a feed's items comes back in 10ms rather than 130.
|
||||||
|
- On Postgres, a write no longer waits for the database server's disk. A scan read about one feed a second and now reads 1,500 in 20 seconds.
|
||||||
|
- Checking every feed no longer makes every icon in the feed list flash while it runs: the list keeps the icons it has already drawn.
|
||||||
|
|
||||||
|
### Security
|
||||||
|
|
||||||
|
- A flood of sign-in attempts no longer stalls the site for everyone else. Passwords are checked a few at a time, away from the threads that answer every other request; forty wrong passwords at a time made everything else take four seconds.
|
||||||
|
- A wrong password is refused in the same time whether or not the name is an account here, so how long it takes no longer tells anyone which names are.
|
||||||
|
|
||||||
## [0.10.0] - 2026-10-05
|
## [0.10.0] - 2026-10-05
|
||||||
|
|
||||||
### Added
|
### Added
|
||||||
@@ -735,7 +758,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
|||||||
- Torrent enclosures through librqbit, seeding to a ratio or a time, with a stall timeout.
|
- Torrent enclosures through librqbit, seeding to a ratio or a time, with a stall timeout.
|
||||||
- `ipx import` and `ipx export` for OPML, and systemd units in `contrib/`.
|
- `ipx import` and `ipx export` for OPML, and systemd units in `contrib/`.
|
||||||
|
|
||||||
[unreleased]: https://git.sdf1.net/rays/ipx/compare/v0.10.0...main
|
[unreleased]: https://git.sdf1.net/rays/ipx/compare/v0.10.2...main
|
||||||
|
[0.10.2]: https://git.sdf1.net/rays/ipx/compare/v0.10.1...v0.10.2
|
||||||
|
[0.10.1]: https://git.sdf1.net/rays/ipx/compare/v0.10.0...v0.10.1
|
||||||
[0.10.0]: https://git.sdf1.net/rays/ipx/compare/v0.9.1...v0.10.0
|
[0.10.0]: https://git.sdf1.net/rays/ipx/compare/v0.9.1...v0.10.0
|
||||||
[0.9.1]: https://git.sdf1.net/rays/ipx/compare/v0.9.0...v0.9.1
|
[0.9.1]: https://git.sdf1.net/rays/ipx/compare/v0.9.0...v0.9.1
|
||||||
[0.9.0]: https://git.sdf1.net/rays/ipx/compare/v0.8.4...v0.9.0
|
[0.9.0]: https://git.sdf1.net/rays/ipx/compare/v0.8.4...v0.9.0
|
||||||
|
|||||||
28
CLAUDE.md
28
CLAUDE.md
@@ -124,7 +124,8 @@ npx tsc -p . # type-checks web/src
|
|||||||
node tests/page-smoke.js
|
node tests/page-smoke.js
|
||||||
node tests/native-bridge.js # the page hands playback to a native shell
|
node tests/native-bridge.js # the page hands playback to a native shell
|
||||||
node tests/contrast.js # every theme's palette against WCAG AA
|
node tests/contrast.js # every theme's palette against WCAG AA
|
||||||
npx playwright test # 40 browser tests against a real daemon on fixture feeds
|
npx playwright test # 66 browser tests against a real daemon on fixture feeds
|
||||||
|
node tests/load/run.js # k6: many people at once against a scratch daemon with 1,500 feeds
|
||||||
```
|
```
|
||||||
|
|
||||||
Things about the browser suite that have cost time:
|
Things about the browser suite that have cost time:
|
||||||
@@ -141,6 +142,31 @@ Things about the browser suite that have cost time:
|
|||||||
* `webServer` starts **before** `globalSetup`, which is why the fixture config is written at
|
* `webServer` starts **before** `globalSetup`, which is why the fixture config is written at
|
||||||
config-load time instead.
|
config-load time instead.
|
||||||
|
|
||||||
|
The load tests (`tests/load`, issue #133) are for what one browser cannot show: twenty-five people
|
||||||
|
browsing at once, a hundred players saving positions while a scan writes, fifty listeners
|
||||||
|
seeking through files, a flood of wrong passwords. Things about them:
|
||||||
|
|
||||||
|
* They need **k6**, which `/src/install.sh` installs from k6's own signed apt repository, and they
|
||||||
|
build and run a **release** binary: Argon2 in a debug build takes about a second a sign-in,
|
||||||
|
which would measure the build.
|
||||||
|
* `run.js` serves the feeds itself, generated, and starts **its own daemon** under
|
||||||
|
`/tmp/ipx-load`, wiped each run, on ports 8793 and 8794, so it can run beside the browser suite.
|
||||||
|
It stops that daemon by its PID. People sign in by the `X-Load-User` header, trusted from
|
||||||
|
127.0.0.1, as production trusts Cloudflare Access's; only the sign-in test uses a password.
|
||||||
|
* Each threshold is a budget well above what the request takes now: it is there to catch a query
|
||||||
|
that has started asking once per feed, or writers queueing on a lock, not a slow minute.
|
||||||
|
* `ipx status`, the Docker healthcheck, runs every second throughout, and a run fails if one takes
|
||||||
|
the healthcheck's 5s.
|
||||||
|
* One at a time: `node tests/load/run.js signin`. All four take about six minutes.
|
||||||
|
* The daemon is on **SQLite** unless `--postgres`, which puts it in `IPX_TEST_DATABASE_URL`, in
|
||||||
|
a schema of its own (`ipx_load`) dropped and made again each run. That URL, in `/src/.envrc`, is
|
||||||
|
`ipodderx_test` on the `databases` project's server, production's, as `ipodderx`, taken from
|
||||||
|
`ipx.env`; the harness refuses a database whose name does not end in `_test`. Postgres's
|
||||||
|
numbers are what production would see, and a run loads the server production's database is on,
|
||||||
|
so not while someone is using iPX in earnest. Each budget has a value for each database.
|
||||||
|
* The Rust tests run on Postgres with the same URL: `. /src/.envrc && cargo test`. Each test
|
||||||
|
makes a schema `ipxt_<pid>_<n>`, and the next run drops those an earlier one left.
|
||||||
|
|
||||||
Non-trivial logic leaves one runnable check behind. Pure functions (`merge_policy`, `pick`,
|
Non-trivial logic leaves one runnable check behind. Pure functions (`merge_policy`, `pick`,
|
||||||
`matches_keywords`, `parse_interval`) are the easiest place to put it.
|
`matches_keywords`, `parse_interval`) are the easiest place to put it.
|
||||||
|
|
||||||
|
|||||||
2
Cargo.lock
generated
2
Cargo.lock
generated
@@ -1820,7 +1820,7 @@ checksum = "791930b43c0d5973160d90a8f3894509f2b273430f5c5c73b668636d0287c5c0"
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "ipx"
|
name = "ipx"
|
||||||
version = "0.10.0"
|
version = "0.10.2"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"ammonia",
|
"ammonia",
|
||||||
"anyhow",
|
"anyhow",
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "ipx"
|
name = "ipx"
|
||||||
version = "0.10.0"
|
version = "0.10.2"
|
||||||
edition = "2024"
|
edition = "2024"
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
|
|||||||
@@ -38,7 +38,9 @@ build needs node and `npm ci` run once.
|
|||||||
re-derived into the database each scan, never written to config.toml. A Patreon creator link
|
re-derived into the database each scan, never written to config.toml. A Patreon creator link
|
||||||
(a token, no `show=`) with more than one show is treated the same way, before any fetch: its
|
(a token, no `show=`) with more than one show is treated the same way, before any fetch: its
|
||||||
shows come from Patreon's web API and each becomes a derived feed.
|
shows come from Patreon's web API and each becomes a derived feed.
|
||||||
4. Record entries. A changed title or description flips the item back to unread.
|
4. Record entries: the new ones, and those the feed now gives a different title, text, artwork,
|
||||||
|
length or number, compared with what is stored for the items it lists (#141). A corrected item
|
||||||
|
keeps its read state.
|
||||||
5. Record enclosures. `enclosures.url` is `UNIQUE`, which is the dedupe key and subsumes the
|
5. Record enclosures. `enclosures.url` is `UNIQUE`, which is the dedupe key and subsumes the
|
||||||
original's `history.dat` pickle: a reaped file keeps its row so it is never fetched twice.
|
original's `history.dat` pickle: a reaped file keeps its row so it is never fetched twice.
|
||||||
6. Apply the merged policy (see [users.md](users.md)) and mark anything rejected as `skipped` with
|
6. Apply the merged policy (see [users.md](users.md)) and mark anything rejected as `skipped` with
|
||||||
@@ -90,13 +92,15 @@ printf '{"cmd":"fetch","force":true}\n' | socat - UNIX-CONNECT:$XDG_RUNTIME_DIR/
|
|||||||
```
|
```
|
||||||
|
|
||||||
**Commands** — `fetch` (optional `feed`, `force`), `reap` (optional `dry_run`), `download`
|
**Commands** — `fetch` (optional `feed`, `force`), `reap` (optional `dry_run`), `download`
|
||||||
(`enclosure`), `status`.
|
(`enclosure`), `status`, `reload` (read the catalogue and settings again from the database, which
|
||||||
|
`ipx add`, `rm` and `import` send after changing them in a process of their own; answered with
|
||||||
|
`status`).
|
||||||
|
|
||||||
**Events** — `feed_start`, `feed_skip`, `feed_done`, `feed_error`, `progress`, `download_done`,
|
**Events** — `feed_start`, `feed_skip`, `feed_done`, `feed_error`, `progress`, `download_done`,
|
||||||
`download_error`, `torrent_deferred`, `reaped`, `reap_done`, `scan_done`, `status`, `error`.
|
`download_error`, `torrent_deferred`, `reaped`, `reap_done`, `scan_done`, `status`, `error`.
|
||||||
`scan_done`, `reap_done` and `status` are terminal: a client that asked for work stops reading
|
`scan_done`, `reap_done` and `status` are terminal: a client that asked for work stops reading
|
||||||
there. Commands run one at a time, in the order they arrive, except `status`: the socket answers it
|
there. Commands run one at a time, in the order they arrive, except `status` and `reload`: the
|
||||||
straight away, so the Docker healthcheck is never left waiting behind a scan or a download, and
|
socket answers them straight away, so the Docker healthcheck is never left waiting behind a scan or a download, and
|
||||||
answers only the client that asked, since `status` would end any other client's session.
|
answers only the client that asked, since `status` would end any other client's session.
|
||||||
|
|
||||||
Progress carries the enclosure id, without which a UI cannot tell one download from another and
|
Progress carries the enclosure id, without which a UI cannot tell one download from another and
|
||||||
@@ -128,7 +132,7 @@ else a `401`. A feed's items and files (its entries, `download-latest`, `/api/en
|
|||||||
| `POST /api/enclosures/{id}/download`, `DELETE /api/enclosures/{id}` | `?force=true` overrides the shared-file warning |
|
| `POST /api/enclosures/{id}/download`, `DELETE /api/enclosures/{id}` | `?force=true` overrides the shared-file warning |
|
||||||
| `POST /api/fetch` | |
|
| `POST /api/fetch` | |
|
||||||
| `GET /api/opml`, `POST /api/opml` | export your subscriptions; subscribe to every feed in an OPML |
|
| `GET /api/opml`, `POST /api/opml` | export your subscriptions; subscribe to every feed in an OPML |
|
||||||
| `GET /api/directory/{id}` | a listed feed's description and latest twenty items, for its page before you subscribe: title, link, date, length and text, never a file or its address, and only for a feed the directory lists |
|
| `GET /api/directory/{id}` | a listed feed's description, why its last check failed when it did (in words, never its error, which can name its address), and latest twenty items, for its page before you subscribe: title, link, date, length and text, never a file or its address, and only for a feed the directory lists |
|
||||||
| `GET /api/popular`, `GET /api/directory`, `POST /api/popular/{id}` | the ten most subscribed feeds, and every listable feed A to Z, with an OPML's feeds in place of the OPML and everyone counted (id, title, art, count, whether it is yours, the feed's iTunes category, whether it carries audio or video; never a URL, never a private feed); subscribe by id |
|
| `GET /api/popular`, `GET /api/directory`, `POST /api/popular/{id}` | the ten most subscribed feeds, and every listable feed A to Z, with an OPML's feeds in place of the OPML and everyone counted (id, title, art, count, whether it is yours, the feed's iTunes category, whether it carries audio or video; never a URL, never a private feed); subscribe by id |
|
||||||
| `GET /api/settings`, `PATCH /api/settings` | admin-only to write |
|
| `GET /api/settings`, `PATCH /api/settings` | admin-only to write |
|
||||||
| `GET /api/users`, `POST /api/users`, `PATCH /api/users/{id}`, `DELETE /api/users/{id}` | admin-only; the only admin cannot be demoted or removed |
|
| `GET /api/users`, `POST /api/users`, `PATCH /api/users/{id}`, `DELETE /api/users/{id}` | admin-only; the only admin cannot be demoted or removed |
|
||||||
@@ -146,6 +150,7 @@ cargo test # parsing, filters, retention, schedules, SQL, per-use
|
|||||||
node tests/page-smoke.js # the page script loads and every selector it wires at load exists
|
node tests/page-smoke.js # the page script loads and every selector it wires at load exists
|
||||||
node tests/native-bridge.js # the page hands playback to a native shell
|
node tests/native-bridge.js # the page hands playback to a native shell
|
||||||
npx playwright test # a real browser against a real daemon on fixture feeds
|
npx playwright test # a real browser against a real daemon on fixture feeds
|
||||||
|
node tests/load/run.js # k6: many people at once against a daemon with 1,500 generated feeds
|
||||||
```
|
```
|
||||||
|
|
||||||
The Rust tests cannot see a wrong selector, a handler that runs and does nothing, or a page that
|
The Rust tests cannot see a wrong selector, a handler that runs and does nothing, or a page that
|
||||||
|
|||||||
@@ -13,13 +13,17 @@ processes must never download the same thing. `--local` forces the work to happe
|
|||||||
| `ipx list` | Subscriptions and their state |
|
| `ipx list` | Subscriptions and their state |
|
||||||
| `ipx status` | Counts: feeds, pending, downloaded |
|
| `ipx status` | Counts: feeds, pending, downloaded |
|
||||||
| `ipx fetch [FEED] [--force]` | Scan everything, or one feed. `--force` ignores the TTL |
|
| `ipx fetch [FEED] [--force]` | Scan everything, or one feed. `--force` ignores the TTL |
|
||||||
| `ipx add <url> [--folder X] [--keywords a,b]` | Subscribe; the id comes from the feed title |
|
| `ipx add <url> [--folder X] [--keywords a,b] [--list] [--category C]` | Subscribe; the id comes from the feed title. `--list` puts it in the Directory for anyone to subscribe to instead, and keeps it there when its last subscriber leaves; `--category` files it under one of the Directory's categories. Run for a feed already in the catalogue, these list it or set its category |
|
||||||
| `ipx rm <feed>` | Unsubscribe; downloads and history are kept |
|
| `ipx rm <feed>` | Unsubscribe; downloads and history are kept |
|
||||||
| `ipx import <file.opml>` / `ipx export <file.opml>` | Move subscriptions in or out. Import subscribes the first admin, as the shared web token does; in the web UI it subscribes whoever is signed in |
|
| `ipx import <file.opml>` / `ipx export <file.opml>` | Move subscriptions in or out. Import subscribes the first admin, as the shared web token does; in the web UI it subscribes whoever is signed in |
|
||||||
| `ipx reap [--dry-run]` | Run retention now |
|
| `ipx reap [--dry-run]` | Run retention now |
|
||||||
| `ipx user <add\|list\|passwd\|rm>` | Accounts for the web UI |
|
| `ipx user <add\|list\|passwd\|rm>` | Accounts for the web UI |
|
||||||
| `ipx daemon [--web ADDR]` | Scheduler, control socket and web UI |
|
| `ipx daemon [--web ADDR]` | Scheduler, control socket and web UI |
|
||||||
|
|
||||||
|
`add`, `rm` and `import` change the catalogue in a process of their own. With a daemon running they
|
||||||
|
then tell it to read the catalogue again; without that it kept its own copy and wrote it back at
|
||||||
|
its next change, undoing them. If it does not answer they say so: restart it.
|
||||||
|
|
||||||
## Accounts
|
## Accounts
|
||||||
|
|
||||||
Passwords are read from **stdin**, so they miss the shell history and any `ps` listing.
|
Passwords are read from **stdin**, so they miss the shell history and any `ps` listing.
|
||||||
|
|||||||
@@ -27,6 +27,12 @@ config.toml's default location is `$XDG_CONFIG_HOME/ipx/config.toml`
|
|||||||
Postgres database instead. Back SQLite up by copying `state.db` while the daemon is stopped, or
|
Postgres database instead. Back SQLite up by copying `state.db` while the daemon is stopped, or
|
||||||
with `sqlite3 state.db .backup`; back Postgres up with `pg_dump`.
|
with `sqlite3 state.db .backup`; back Postgres up with `pg_dump`.
|
||||||
|
|
||||||
|
ipx does not wait for the disk to confirm each write: SQLite syncs at checkpoints
|
||||||
|
(`synchronous=NORMAL`), and ipx's own Postgres sessions set `synchronous_commit=off`, which leaves
|
||||||
|
other databases on the server as they are. Waiting made a scan take most of a second a feed. A
|
||||||
|
crash of ipx loses nothing; a power cut, or a crash of the Postgres server, can lose the last moment
|
||||||
|
of changes, a read mark or a saved position, and never corrupts the database.
|
||||||
|
|
||||||
## `[general]`
|
## `[general]`
|
||||||
|
|
||||||
```toml
|
```toml
|
||||||
|
|||||||
@@ -87,6 +87,10 @@ Patreon or Supercast, which put the key in the path, and any feed inside an OPML
|
|||||||
itself. Those are someone's paid subscriptions, and listing them would let anyone here read what
|
itself. Those are someone's paid subscriptions, and listing them would let anyone here read what
|
||||||
they pay for.
|
they pay for.
|
||||||
|
|
||||||
|
**Currently Listening**, under the Directory, is every episode you started and have not finished,
|
||||||
|
across all your feeds, with how much is left. A click picks one up where you left off. The search
|
||||||
|
box looks through it, as it does a feed's items, and through the Directory by name.
|
||||||
|
|
||||||
An admin can do the same from **Settings → Manage users…**: add someone (with a password, or none
|
An admin can do the same from **Settings → Manage users…**: add someone (with a password, or none
|
||||||
for someone the proxy signs in), tick or untick Admin, or remove an account. Removing one takes its
|
for someone the proxy signs in), tick or untick Admin, or remove an account. Removing one takes its
|
||||||
subscriptions and read state with it; downloaded files stay. The only admin cannot be demoted or
|
subscriptions and read state with it; downloaded files stay. The only admin cannot be demoted or
|
||||||
|
|||||||
@@ -7,7 +7,8 @@
|
|||||||
"typecheck": "tsc -p .",
|
"typecheck": "tsc -p .",
|
||||||
"smoke": "node tests/page-smoke.js",
|
"smoke": "node tests/page-smoke.js",
|
||||||
"test": "playwright test",
|
"test": "playwright test",
|
||||||
"test:headed": "playwright test --headed"
|
"test:headed": "playwright test --headed",
|
||||||
|
"load": "node tests/load/run.js"
|
||||||
},
|
},
|
||||||
"devDependencies": {
|
"devDependencies": {
|
||||||
"@playwright/test": "^1.56.0",
|
"@playwright/test": "^1.56.0",
|
||||||
|
|||||||
39
src/auth.rs
39
src/auth.rs
@@ -30,6 +30,35 @@ pub fn verify_password(password: &str, stored: &str) -> bool {
|
|||||||
.is_ok()
|
.is_ok()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// How many password checks run at once: each is tens of milliseconds of CPU, and forty wrong
|
||||||
|
/// passwords at a time, run on the async workers, made every other request wait 4s (#137).
|
||||||
|
/// Half the cores, so a flood of sign-ins waits on itself and the rest of the server has the rest.
|
||||||
|
static CHECKS: std::sync::LazyLock<tokio::sync::Semaphore> = std::sync::LazyLock::new(|| {
|
||||||
|
tokio::sync::Semaphore::new(std::thread::available_parallelism().map_or(1, |n| (n.get() / 2).max(1)))
|
||||||
|
});
|
||||||
|
|
||||||
|
/// What a name that is not an account is checked against, so that refusing it takes as long as
|
||||||
|
/// refusing a wrong password does. Refused without a check, it came back 31ms sooner, and the
|
||||||
|
/// time told anyone which names are accounts here (#138).
|
||||||
|
static DECOY: std::sync::LazyLock<String> =
|
||||||
|
std::sync::LazyLock::new(|| hash_password("no account here has this password").expect("hashing a fixed password"));
|
||||||
|
|
||||||
|
/// A sign-in's password against the account's hash, or against `DECOY` when there is no account
|
||||||
|
/// or it has no password: false either way, in the same time. Off the async workers, and a few at
|
||||||
|
/// a time; see `CHECKS`.
|
||||||
|
pub async fn check_password(password: String, stored: Option<String>) -> bool {
|
||||||
|
let Ok(_turn) = CHECKS.acquire().await else { return false };
|
||||||
|
tokio::task::spawn_blocking(move || match stored {
|
||||||
|
Some(h) => verify_password(&password, &h),
|
||||||
|
None => {
|
||||||
|
verify_password(&password, &DECOY);
|
||||||
|
false
|
||||||
|
}
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.unwrap_or(false)
|
||||||
|
}
|
||||||
|
|
||||||
/// A session id: 256 bits of urandom, hex. Long enough that guessing is not a strategy.
|
/// A session id: 256 bits of urandom, hex. Long enough that guessing is not a strategy.
|
||||||
pub fn new_session_token() -> String {
|
pub fn new_session_token() -> String {
|
||||||
let mut bytes = [0u8; 32];
|
let mut bytes = [0u8; 32];
|
||||||
@@ -76,6 +105,16 @@ pub fn name_from_header(raw: &str) -> Option<String> {
|
|||||||
mod tests {
|
mod tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn a_sign_in_without_an_account_is_checked_all_the_same() {
|
||||||
|
let h = hash_password("correct horse battery").unwrap();
|
||||||
|
assert!(check_password("correct horse battery".into(), Some(h.clone())).await);
|
||||||
|
assert!(!check_password("wrong".into(), Some(h)).await);
|
||||||
|
assert!(!check_password("correct horse battery".into(), None).await);
|
||||||
|
// A decoy that does not parse is refused before any hashing, and the time says so (#138).
|
||||||
|
assert!(PasswordHash::new(&DECOY).is_ok());
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn a_password_verifies_only_against_itself() {
|
fn a_password_verifies_only_against_itself() {
|
||||||
let h = hash_password("correct horse battery").unwrap();
|
let h = hash_password("correct horse battery").unwrap();
|
||||||
|
|||||||
134
src/db.rs
134
src/db.rs
@@ -125,6 +125,21 @@ async fn connect(location: &str) -> Result<sea_orm::DatabaseConnection> {
|
|||||||
// 52ms, and 39s to subscribe an admin to 1,500 feeds (#135). A crash of ipx loses nothing;
|
// 52ms, and 39s to subscribe an admin to 1,500 feeds (#135). A crash of ipx loses nothing;
|
||||||
// only a power cut can lose the last transactions, and it never corrupts the file.
|
// only a power cut can lose the last transactions, and it never corrupts the file.
|
||||||
opts.map_sqlx_sqlite_opts(|o| o.synchronous(sea_orm::sqlx::sqlite::SqliteSynchronous::Normal));
|
opts.map_sqlx_sqlite_opts(|o| o.synchronous(sea_orm::sqlx::sqlite::SqliteSynchronous::Normal));
|
||||||
|
// The same on Postgres, for ipx's own sessions: a commit waited for the WAL to reach the
|
||||||
|
// disk, 12.8ms on the databases project's server, and at about fifty commits a feed a scan
|
||||||
|
// took most of a second a feed (#140). A crash of the Postgres server can lose the last
|
||||||
|
// commits, about the last 0.6s, and never corrupts anything. Added to whatever options the
|
||||||
|
// URL sets, so it holds for the tests' search_path too.
|
||||||
|
opts.map_sqlx_postgres_opts(|o| o.options([("synchronous_commit", "off")]));
|
||||||
|
// Several connections, where SeaORM gives SQLite one unless told: every request, every scan
|
||||||
|
// write and every read went through it in turn, and twenty-five people browsing waited in
|
||||||
|
// line for it at 50 requests a second (#136). In WAL readers run beside the one writer. Each
|
||||||
|
// transaction here writes first, so it takes the write lock or waits out the busy timeout for
|
||||||
|
// it; SQLite refuses at once only a transaction that read a snapshot a write has since moved
|
||||||
|
// past. Postgres keeps sqlx's ten.
|
||||||
|
if !is_postgres(location) {
|
||||||
|
opts.max_connections(8);
|
||||||
|
}
|
||||||
sea_orm::Database::connect(opts)
|
sea_orm::Database::connect(opts)
|
||||||
.await
|
.await
|
||||||
.with_context(|| format!("opening {}", redact(location)))
|
.with_context(|| format!("opening {}", redact(location)))
|
||||||
@@ -313,6 +328,8 @@ impl Db {
|
|||||||
}
|
}
|
||||||
let path = std::env::temp_dir().join(format!("ipx-test-{pid}-{n}.db"));
|
let path = std::env::temp_dir().join(format!("ipx-test-{pid}-{n}.db"));
|
||||||
let orm = connect(&path.display().to_string()).await?;
|
let orm = connect(&path.display().to_string()).await?;
|
||||||
|
// As Db::open has it, so the tests' connections share the file as the daemon's do.
|
||||||
|
orm.execute_unprepared("PRAGMA journal_mode = WAL").await?;
|
||||||
create_missing(&orm).await?;
|
create_missing(&orm).await?;
|
||||||
Ok(Self { orm, tmp: Some(path) })
|
Ok(Self { orm, tmp: Some(path) })
|
||||||
}
|
}
|
||||||
@@ -2042,6 +2059,64 @@ impl Db {
|
|||||||
))
|
))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Of these items, the stored ones the feed now tells differently: a title, text, artwork,
|
||||||
|
/// length or number that is not what is stored, a field the feed leaves out counting as
|
||||||
|
/// unchanged, as record_entry writes them. A scan inserted only items it had not stored
|
||||||
|
/// (#96), and with that alone a publisher's correction never arrived: Ain't It Cool News's
|
||||||
|
/// items kept a broken copy, each with the text of the one before it, after the feed was
|
||||||
|
/// fixed (#141). Only what the feed lists is read back, so a scan reads as much as it was
|
||||||
|
/// sent, not the archive.
|
||||||
|
#[tracing::instrument(skip_all)]
|
||||||
|
pub async fn changed_items(
|
||||||
|
&self,
|
||||||
|
feed_id: &str,
|
||||||
|
entries: &[crate::feed::Entry],
|
||||||
|
) -> Result<std::collections::HashSet<String>> {
|
||||||
|
if entries.is_empty() {
|
||||||
|
return Ok(Default::default());
|
||||||
|
}
|
||||||
|
let mut a = Args::default();
|
||||||
|
let feed = a.p(feed_id);
|
||||||
|
let guids: Vec<String> = entries.iter().map(|e| a.p(e.guid.clone())).collect();
|
||||||
|
let sql = format!(
|
||||||
|
"SELECT guid, title, description, image, duration, episode, season FROM entries
|
||||||
|
WHERE feed_id = {feed} AND guid IN ({})",
|
||||||
|
guids.join(", ")
|
||||||
|
);
|
||||||
|
type Stored = (Option<String>, Option<String>, Option<String>, Option<i64>, Option<i64>, Option<i64>);
|
||||||
|
let mut stored = std::collections::HashMap::<String, Stored>::new();
|
||||||
|
for r in self.rows(&sql, a.0).await? {
|
||||||
|
stored.insert(
|
||||||
|
r.try_get("", "guid")?,
|
||||||
|
(
|
||||||
|
r.try_get("", "title")?,
|
||||||
|
r.try_get("", "description")?,
|
||||||
|
r.try_get("", "image")?,
|
||||||
|
r.try_get("", "duration")?,
|
||||||
|
r.try_get("", "episode")?,
|
||||||
|
r.try_get("", "season")?,
|
||||||
|
),
|
||||||
|
);
|
||||||
|
}
|
||||||
|
fn differs<T: PartialEq>(new: &Option<T>, old: &Option<T>) -> bool {
|
||||||
|
new.is_some() && new != old
|
||||||
|
}
|
||||||
|
Ok(entries
|
||||||
|
.iter()
|
||||||
|
.filter(|e| {
|
||||||
|
stored.get(&e.guid).is_some_and(|s| {
|
||||||
|
differs(&e.title, &s.0)
|
||||||
|
|| differs(&e.description, &s.1)
|
||||||
|
|| differs(&e.image, &s.2)
|
||||||
|
|| differs(&e.duration, &s.3)
|
||||||
|
|| differs(&e.episode, &s.4)
|
||||||
|
|| differs(&e.season, &s.5)
|
||||||
|
})
|
||||||
|
})
|
||||||
|
.map(|e| e.guid.clone())
|
||||||
|
.collect())
|
||||||
|
}
|
||||||
|
|
||||||
#[tracing::instrument(skip_all)]
|
#[tracing::instrument(skip_all)]
|
||||||
pub async fn skipped_by_filter(&self, feed_id: &str) -> Result<std::collections::HashMap<String, String>> {
|
pub async fn skipped_by_filter(&self, feed_id: &str) -> Result<std::collections::HashMap<String, String>> {
|
||||||
self.rows(
|
self.rows(
|
||||||
@@ -2106,13 +2181,37 @@ mod tests {
|
|||||||
use super::*;
|
use super::*;
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn sqlite_syncs_at_checkpoints_not_every_commit() {
|
async fn a_read_does_not_wait_for_a_write_in_progress() {
|
||||||
|
use sea_orm::TransactionTrait;
|
||||||
let db = Db::memory().await.unwrap();
|
let db = Db::memory().await.unwrap();
|
||||||
if db.orm.get_database_backend() != sea_orm::DbBackend::Sqlite {
|
// A write held open, as a scan's or an import's is while it runs.
|
||||||
return;
|
let tx = db.orm.begin().await.unwrap();
|
||||||
|
tx.execute_unprepared("INSERT INTO feeds (id, url) VALUES ('held', 'http://x/held')").await.unwrap();
|
||||||
|
// With one connection, which SQLite had (#136), this waited for the transaction to end.
|
||||||
|
let read = tokio::time::timeout(std::time::Duration::from_secs(2), db.feed_summary("other")).await;
|
||||||
|
assert!(read.is_ok(), "a read waited behind a write");
|
||||||
|
tx.commit().await.unwrap();
|
||||||
|
// And writers at once all get their turn, not a refusal.
|
||||||
|
let writes = (0..16).map(|i| {
|
||||||
|
let db = &db;
|
||||||
|
async move { db.record_feed(&format!("f{i}"), "http://x/f", Some("F"), None, None, None, None, None, None).await }
|
||||||
|
});
|
||||||
|
for r in futures_util::future::join_all(writes).await {
|
||||||
|
r.unwrap();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn a_commit_does_not_wait_for_the_disk() {
|
||||||
|
let db = Db::memory().await.unwrap();
|
||||||
|
let backend = db.orm.get_database_backend();
|
||||||
|
if backend == sea_orm::DbBackend::Sqlite {
|
||||||
|
let row = db.orm.query_one_raw(Statement::from_string(backend, "PRAGMA synchronous")).await.unwrap();
|
||||||
|
assert_eq!(row.unwrap().try_get_by_index::<i32>(0).unwrap(), 1, "NORMAL, on every connection in the pool (#135)");
|
||||||
|
} else {
|
||||||
|
let row = db.orm.query_one_raw(Statement::from_string(backend, "SHOW synchronous_commit")).await.unwrap();
|
||||||
|
assert_eq!(row.unwrap().try_get_by_index::<String>(0).unwrap(), "off", "with the test's own search_path as well (#140)");
|
||||||
}
|
}
|
||||||
let row = db.orm.query_one_raw(Statement::from_string(sea_orm::DbBackend::Sqlite, "PRAGMA synchronous")).await.unwrap();
|
|
||||||
assert_eq!(row.unwrap().try_get_by_index::<i32>(0).unwrap(), 1, "NORMAL, on every connection in the pool (#135)");
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
@@ -2355,6 +2454,31 @@ mod tests {
|
|||||||
assert_eq!(files, ["u1"].map(String::from).into()); // u2 is g's: left to the insert to find
|
assert_eq!(files, ["u1"].map(String::from).into()); // u2 is g's: left to the insert to find
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn a_scan_finds_the_stored_items_the_feed_now_tells_differently() {
|
||||||
|
let db = Db::memory().await.unwrap();
|
||||||
|
let item = |guid: &str, title: &str, text: Option<&str>| crate::feed::Entry {
|
||||||
|
guid: guid.into(),
|
||||||
|
title: Some(title.into()),
|
||||||
|
description: text.map(str::to_owned),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
db.record_entry("f", &item("a", "A", Some("the review"))).await.unwrap();
|
||||||
|
db.record_entry("f", &item("b", "B", Some("the next review"))).await.unwrap();
|
||||||
|
db.record_entry("f", &item("c", "C", Some("as it was"))).await.unwrap();
|
||||||
|
let changed = db
|
||||||
|
.changed_items("f", &[
|
||||||
|
item("a", "A", Some("the review, fixed")), // the text corrected
|
||||||
|
item("b", "B renamed", Some("the next review")), // retitled
|
||||||
|
item("c", "C", None), // says nothing of its text: what is stored stands
|
||||||
|
item("d", "D", Some("new")), // not stored: the insert's, not this
|
||||||
|
])
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(changed, ["a", "b"].map(String::from).into());
|
||||||
|
assert!(db.changed_items("f", &[]).await.unwrap().is_empty());
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn only_artwork_a_feed_names_is_fetched_for_the_page() {
|
async fn only_artwork_a_feed_names_is_fetched_for_the_page() {
|
||||||
let db = Db::memory().await.unwrap();
|
let db = Db::memory().await.unwrap();
|
||||||
|
|||||||
@@ -190,6 +190,9 @@ pub fn explain_failure(msg: &str) -> Option<Failure> {
|
|||||||
if low.contains("http 401") || low.contains("http 403") {
|
if low.contains("http 401") || low.contains("http 403") {
|
||||||
return Some(Failure { reason: "The site refuses ipx's requests.", new_url: None });
|
return Some(Failure { reason: "The site refuses ipx's requests.", new_url: None });
|
||||||
}
|
}
|
||||||
|
if low.contains("answered with nothing") {
|
||||||
|
return Some(Failure { reason: "This address answers with nothing; the site may be gone.", new_url: None });
|
||||||
|
}
|
||||||
if low.contains("http 402") {
|
if low.contains("http 402") {
|
||||||
return Some(Failure { reason: "The feed now needs a paid plan.", new_url: None });
|
return Some(Failure { reason: "The feed now needs a paid plan.", new_url: None });
|
||||||
}
|
}
|
||||||
@@ -1310,6 +1313,10 @@ mod tests {
|
|||||||
assert_eq!(explain_failure("HTTP 401 Unauthorized").unwrap().reason, "The site refuses ipx's requests.");
|
assert_eq!(explain_failure("HTTP 401 Unauthorized").unwrap().reason, "The site refuses ipx's requests.");
|
||||||
assert_eq!(explain_failure("HTTP 403 Forbidden").unwrap().reason, "The site refuses ipx's requests.");
|
assert_eq!(explain_failure("HTTP 403 Forbidden").unwrap().reason, "The site refuses ipx's requests.");
|
||||||
assert_eq!(explain_failure("HTTP 402 Payment Required").unwrap().reason, "The feed now needs a paid plan.");
|
assert_eq!(explain_failure("HTTP 402 Payment Required").unwrap().reason, "The feed now needs a paid plan.");
|
||||||
|
assert_eq!(
|
||||||
|
explain_failure("the feed's address answered with nothing at all").unwrap().reason,
|
||||||
|
"This address answers with nothing; the site may be gone."
|
||||||
|
);
|
||||||
let dns = explain_failure("connecting: dns error: failed to lookup address information").unwrap();
|
let dns = explain_failure("connecting: dns error: failed to lookup address information").unwrap();
|
||||||
assert_eq!(dns.reason, "This address no longer resolves; the site is gone.");
|
assert_eq!(dns.reason, "This address no longer resolves; the site is gone.");
|
||||||
let moved = explain_failure("got a web page, not a feed; it links https://x/feed as its feed").unwrap();
|
let moved = explain_failure("got a web page, not a feed; it links https://x/feed as its feed").unwrap();
|
||||||
|
|||||||
45
src/ipc.rs
45
src/ipc.rs
@@ -109,6 +109,11 @@ pub enum Command {
|
|||||||
enclosure: i64,
|
enclosure: i64,
|
||||||
},
|
},
|
||||||
Status,
|
Status,
|
||||||
|
/// Read the catalogue and server settings again from the database, which `ipx add`, `rm`
|
||||||
|
/// and `import` change in a process of their own. A daemon keeps them in memory, and with
|
||||||
|
/// its old copy wrote a removed feed back, and dropped an added one, at its next change of
|
||||||
|
/// its own (#120). Answered, as status is, with the counts after it.
|
||||||
|
Reload,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Where events go: the socket, the terminal, or both.
|
/// Where events go: the socket, the terminal, or both.
|
||||||
@@ -208,13 +213,15 @@ pub async fn daemon_is_live(path: &Path) -> bool {
|
|||||||
UnixStream::connect(path).await.is_ok()
|
UnixStream::connect(path).await.is_ok()
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Answers `status` for the socket, without the worker. The worker runs one job at a time, and a
|
/// Answers `status` and `reload` for the socket, without the worker. The worker runs one job at a
|
||||||
/// healthcheck left waiting behind a scan or a long download timed out and called a busy daemon
|
/// time, and a healthcheck left waiting behind a scan or a long download timed out and called a
|
||||||
/// dead. The answer goes to the client that asked and no one else: broadcast, it ended any
|
/// busy daemon dead; a reload left waiting would let the web UI write the old catalogue back in
|
||||||
/// `ipx fetch` that was watching a scan, since `status` is a terminal event.
|
/// the meantime. The answer goes to the client that asked and no one else: broadcast, it ended
|
||||||
|
/// any `ipx fetch` that was watching a scan, since `status` is a terminal event.
|
||||||
/// A future, since reading the counts is a database query.
|
/// A future, since reading the counts is a database query.
|
||||||
pub type StatusFn =
|
pub type StatusFn = std::sync::Arc<
|
||||||
std::sync::Arc<dyn Fn() -> std::pin::Pin<Box<dyn std::future::Future<Output = Event> + Send>> + Send + Sync>;
|
dyn Fn(Command) -> std::pin::Pin<Box<dyn std::future::Future<Output = Event> + Send>> + Send + Sync,
|
||||||
|
>;
|
||||||
|
|
||||||
/// Accepts connections, feeding commands to `cmds` and events from `events` back out.
|
/// Accepts connections, feeding commands to `cmds` and events from `events` back out.
|
||||||
pub async fn serve(
|
pub async fn serve(
|
||||||
@@ -285,9 +292,9 @@ async fn handle(
|
|||||||
}
|
}
|
||||||
match serde_json::from_str::<Command>(line) {
|
match serde_json::from_str::<Command>(line) {
|
||||||
// Answered here, not queued behind whatever the worker is on: see StatusFn.
|
// Answered here, not queued behind whatever the worker is on: see StatusFn.
|
||||||
Ok(Command::Status) => {
|
Ok(cmd @ (Command::Status | Command::Reload)) => {
|
||||||
tracing::info!(target: "ipx::io", "-> {line}");
|
tracing::info!(target: "ipx::io", "-> {line}");
|
||||||
let ev = status().await;
|
let ev = status(cmd).await;
|
||||||
log_event(&ev, true);
|
log_event(&ev, true);
|
||||||
let _ = reply.send(ev).await;
|
let _ = reply.send(ev).await;
|
||||||
}
|
}
|
||||||
@@ -303,6 +310,26 @@ async fn handle(
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Tells a running daemon to read the catalogue again (see Command::Reload), and waits for it.
|
||||||
|
pub async fn reload(path: &Path) -> Result<()> {
|
||||||
|
let stream = UnixStream::connect(path).await?;
|
||||||
|
let (read, mut write) = stream.into_split();
|
||||||
|
write.write_all(b"{\"cmd\":\"reload\"}\n").await?;
|
||||||
|
let mut lines = BufReader::new(read).lines();
|
||||||
|
let answer = async {
|
||||||
|
// Other clients' events come here too; the answer is the status, or an error.
|
||||||
|
while let Some(line) = lines.next_line().await? {
|
||||||
|
match serde_json::from_str::<Event>(&line) {
|
||||||
|
Ok(Event::Status { .. }) => return Ok(()),
|
||||||
|
Ok(Event::Error { msg }) => anyhow::bail!("{msg}"),
|
||||||
|
_ => continue,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
anyhow::bail!("the daemon closed the connection without answering")
|
||||||
|
};
|
||||||
|
tokio::time::timeout(std::time::Duration::from_secs(10), answer).await.context("the daemon did not answer in 10s")?
|
||||||
|
}
|
||||||
|
|
||||||
/// Sends one command to a running daemon and prints the events it produces.
|
/// Sends one command to a running daemon and prints the events it produces.
|
||||||
pub async fn proxy(path: &Path, cmd: &Command) -> Result<()> {
|
pub async fn proxy(path: &Path, cmd: &Command) -> Result<()> {
|
||||||
let stream = UnixStream::connect(path).await?;
|
let stream = UnixStream::connect(path).await?;
|
||||||
@@ -411,7 +438,7 @@ mod tests {
|
|||||||
// would end its session.
|
// would end its session.
|
||||||
let mut watcher = events.subscribe();
|
let mut watcher = events.subscribe();
|
||||||
let status: StatusFn =
|
let status: StatusFn =
|
||||||
std::sync::Arc::new(|| Box::pin(async { Event::Status { feeds: 1, pending: 2, downloaded: 3 } }));
|
std::sync::Arc::new(|_| Box::pin(async { Event::Status { feeds: 1, pending: 2, downloaded: 3 } }));
|
||||||
let (client, server) = UnixStream::pair().unwrap();
|
let (client, server) = UnixStream::pair().unwrap();
|
||||||
tokio::spawn(handle(server, events.subscribe(), cmds, status));
|
tokio::spawn(handle(server, events.subscribe(), cmds, status));
|
||||||
|
|
||||||
|
|||||||
121
src/main.rs
121
src/main.rs
@@ -249,9 +249,13 @@ async fn main() -> Result<()> {
|
|||||||
return ipc::proxy(&cfg.general.socket, cmd).await;
|
return ipc::proxy(&cfg.general.socket, cmd).await;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
let socket = cfg.general.socket.clone();
|
||||||
let cfg = assemble_config(&db, cfg, &config_path).await?;
|
let cfg = assemble_config(&db, cfg, &config_path).await?;
|
||||||
|
|
||||||
let is_daemon = matches!(cli.command, Command::Daemon { .. });
|
let is_daemon = matches!(cli.command, Command::Daemon { .. });
|
||||||
|
// These write the catalogue from a process of their own, which a running daemon, holding its
|
||||||
|
// own copy, would write over at its next change (#120).
|
||||||
|
let changes_catalogue = matches!(cli.command, Command::Add { .. } | Command::Rm { .. } | Command::Import { .. });
|
||||||
let (events, _) = broadcast::channel(1024);
|
let (events, _) = broadcast::channel(1024);
|
||||||
let ctx = Arc::new(Ctx {
|
let ctx = Arc::new(Ctx {
|
||||||
cfg: std::sync::RwLock::new(std::sync::Arc::new(cfg)),
|
cfg: std::sync::RwLock::new(std::sync::Arc::new(cfg)),
|
||||||
@@ -285,6 +289,11 @@ async fn main() -> Result<()> {
|
|||||||
Command::Export { file } => export(&ctx, &file).await,
|
Command::Export { file } => export(&ctx, &file).await,
|
||||||
_ => run(&ctx, wire_cmd.expect("only List and Daemon have no wire form")).await,
|
_ => run(&ctx, wire_cmd.expect("only List and Daemon have no wire form")).await,
|
||||||
};
|
};
|
||||||
|
if result.is_ok() && changes_catalogue && ipc::daemon_is_live(&socket).await {
|
||||||
|
if let Err(e) = ipc::reload(&socket).await {
|
||||||
|
eprintln!("the running daemon did not take the change ({e:#}); restart it, or it may undo it");
|
||||||
|
}
|
||||||
|
}
|
||||||
// The batch exporter holds the last few seconds of spans; without this they are lost.
|
// The batch exporter holds the last few seconds of spans; without this they are lost.
|
||||||
if let Some(p) = otel {
|
if let Some(p) = otel {
|
||||||
let _ = p.shutdown();
|
let _ = p.shutdown();
|
||||||
@@ -452,9 +461,19 @@ async fn run(ctx: &Arc<Ctx>, cmd: Cmd) -> Result<()> {
|
|||||||
ctx.out.emit(status(ctx).await);
|
ctx.out.emit(status(ctx).await);
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
Cmd::Reload => reload(ctx).await,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// The catalogue and server settings as the database now has them, in place of the copy in
|
||||||
|
/// memory, for a daemon told another ipx changed them (#120); config.toml, read again, for where
|
||||||
|
/// things are and who gets in.
|
||||||
|
async fn reload(ctx: &Ctx) -> Result<()> {
|
||||||
|
let cfg = assemble_config(&ctx.db, config::Config::load(&ctx.config_path)?, &ctx.config_path).await?;
|
||||||
|
ctx.set_cfg(cfg);
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
/// The counts `ipx status` prints, and /api/status serves. A running daemon's socket answers with
|
/// The counts `ipx status` prints, and /api/status serves. A running daemon's socket answers with
|
||||||
/// this directly rather than through the job queue.
|
/// this directly rather than through the job queue.
|
||||||
pub(crate) async fn status(ctx: &Ctx) -> Event {
|
pub(crate) async fn status(ctx: &Ctx) -> Event {
|
||||||
@@ -511,12 +530,20 @@ async fn daemon(
|
|||||||
let (tx_cmd, mut rx_cmd) = mpsc::channel::<Cmd>(64);
|
let (tx_cmd, mut rx_cmd) = mpsc::channel::<Cmd>(64);
|
||||||
|
|
||||||
let web = start_web(&ctx, &config_path, web_addr, &tx_cmd, &events).await?;
|
let web = start_web(&ctx, &config_path, web_addr, &tx_cmd, &events).await?;
|
||||||
// status is answered by the socket itself; everything else waits its turn in the queue.
|
// status and reload are answered by the socket itself; everything else waits its turn in the
|
||||||
|
// queue.
|
||||||
let answer: ipc::StatusFn = {
|
let answer: ipc::StatusFn = {
|
||||||
let ctx = ctx.clone();
|
let ctx = ctx.clone();
|
||||||
Arc::new(move || {
|
Arc::new(move |cmd| {
|
||||||
let ctx = ctx.clone();
|
let ctx = ctx.clone();
|
||||||
Box::pin(async move { status(&ctx).await })
|
Box::pin(async move {
|
||||||
|
if matches!(cmd, Cmd::Reload)
|
||||||
|
&& let Err(e) = reload(&ctx).await
|
||||||
|
{
|
||||||
|
return Event::Error { msg: format!("reading the catalogue again: {e:#}") };
|
||||||
|
}
|
||||||
|
status(&ctx).await
|
||||||
|
})
|
||||||
})
|
})
|
||||||
};
|
};
|
||||||
let server = tokio::spawn(ipc::serve(socket.clone(), events.clone(), tx_cmd, answer));
|
let server = tokio::spawn(ipc::serve(socket.clone(), events.clone(), tx_cmd, answer));
|
||||||
@@ -1455,6 +1482,13 @@ async fn scan_one(
|
|||||||
};
|
};
|
||||||
|
|
||||||
if bytes.iter().all(u8::is_ascii_whitespace) {
|
if bytes.iter().all(u8::is_ascii_whitespace) {
|
||||||
|
// From a feed that has published, nothing new, said badly (see Outcome::Empty). From one
|
||||||
|
// that has never stored an item, a feed that is not there: a lapsed domain behind a DNS
|
||||||
|
// filter's block page answers 200 and nothing, and looked like a show that had not
|
||||||
|
// posted yet, with no error to act on (#119). The next real read clears it.
|
||||||
|
if stored.entries == 0 {
|
||||||
|
anyhow::bail!("the feed's address answered with nothing at all");
|
||||||
|
}
|
||||||
ctx.db.touch_feed(id, &feed_cfg.url).await?;
|
ctx.db.touch_feed(id, &feed_cfg.url).await?;
|
||||||
return Ok(Outcome::Empty);
|
return Ok(Outcome::Empty);
|
||||||
}
|
}
|
||||||
@@ -1564,13 +1598,17 @@ async fn scan_one(
|
|||||||
// changed nothing however often the feed was scanned.
|
// changed nothing however often the feed was scanned.
|
||||||
let skipped = ctx.db.skipped_by_filter(id).await?;
|
let skipped = ctx.db.skipped_by_filter(id).await?;
|
||||||
let (known_items, known_files) = ctx.db.stored_items(id).await?;
|
let (known_items, known_files) = ctx.db.stored_items(id).await?;
|
||||||
|
let changed = ctx.db.changed_items(id, &parsed.entries).await?;
|
||||||
let mut scan = Scan::default();
|
let mut scan = Scan::default();
|
||||||
// Its own span: the time a feed spends after its fetch was untraced (#96).
|
// Its own span: the time a feed spends after its fetch was untraced (#96).
|
||||||
let store = tracing::info_span!("store", items = parsed.entries.len());
|
let store = tracing::info_span!("store", items = parsed.entries.len());
|
||||||
tracing::Instrument::instrument(async {
|
tracing::Instrument::instrument(async {
|
||||||
for entry in &parsed.entries {
|
for entry in &parsed.entries {
|
||||||
// Only what is not stored yet is inserted; the insert would find the rest and do nothing.
|
// What is not stored yet is inserted, and what the feed has changed since is written
|
||||||
if !known_items.contains(&entry.guid) && ctx.db.record_entry(id, entry).await? {
|
// again (#141); the rest is left alone. A correction keeps the item's read state.
|
||||||
|
if (!known_items.contains(&entry.guid) || changed.contains(&entry.guid))
|
||||||
|
&& ctx.db.record_entry(id, entry).await?
|
||||||
|
{
|
||||||
scan.new_entries += 1;
|
scan.new_entries += 1;
|
||||||
}
|
}
|
||||||
for enc in &entry.enclosures {
|
for enc in &entry.enclosures {
|
||||||
@@ -2248,6 +2286,79 @@ mod tests {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn a_reload_takes_the_catalogue_another_ipx_wrote() {
|
||||||
|
let with = |id: &str| {
|
||||||
|
let mut cfg = config::Config::default();
|
||||||
|
cfg.feeds.insert(id.into(), config::Feed { url: format!("http://x/{id}.xml"), ..feed() });
|
||||||
|
cfg
|
||||||
|
};
|
||||||
|
let held = with("kept-in-memory");
|
||||||
|
let ctx = test_ctx(held.clone()).await;
|
||||||
|
ctx.db.store_config(&held).await.unwrap();
|
||||||
|
// `ipx add` and `ipx rm`, from a process of their own (#120).
|
||||||
|
ctx.db.store_config(&with("added-by-the-cli")).await.unwrap();
|
||||||
|
reload(&ctx).await.unwrap();
|
||||||
|
let urls: Vec<String> = ctx.cfg().feeds.values().map(|f| f.url.clone()).collect();
|
||||||
|
assert_eq!(urls, vec!["http://x/added-by-the-cli.xml".to_string()]);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Answers /empty with a 200 and nothing, /fixed with a feed of one item as it should read,
|
||||||
|
/// anything else with the same item's text broken.
|
||||||
|
async fn empty_or_feed_server() -> String {
|
||||||
|
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||||
|
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||||
|
let addr = listener.local_addr().unwrap();
|
||||||
|
tokio::spawn(async move {
|
||||||
|
loop {
|
||||||
|
let Ok((mut sock, _)) = listener.accept().await else { return };
|
||||||
|
tokio::spawn(async move {
|
||||||
|
let mut buf = [0u8; 2048];
|
||||||
|
let n = sock.read(&mut buf).await.unwrap_or(0);
|
||||||
|
let req = String::from_utf8_lossy(&buf[..n]).into_owned();
|
||||||
|
let body = if req.starts_with("GET /empty") {
|
||||||
|
"\n"
|
||||||
|
} else if req.starts_with("GET /fixed") {
|
||||||
|
"<?xml version=\"1.0\"?><rss version=\"2.0\"><channel><title>T</title><item><title>One</title><guid>g1</guid><description>As it should read.</description></item></channel></rss>"
|
||||||
|
} else {
|
||||||
|
"<?xml version=\"1.0\"?><rss version=\"2.0\"><channel><title>T</title><item><title>One</title><guid>g1</guid><description>Somebody else's text.</description></item></channel></rss>"
|
||||||
|
};
|
||||||
|
let resp = format!("HTTP/1.1 200 OK\r\nContent-Length: {}\r\n\r\n{body}", body.len());
|
||||||
|
let _ = sock.write_all(resp.as_bytes()).await;
|
||||||
|
});
|
||||||
|
}
|
||||||
|
});
|
||||||
|
format!("http://{addr}")
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn a_publishers_correction_reaches_an_item_already_stored() {
|
||||||
|
let base = empty_or_feed_server().await;
|
||||||
|
let ctx = Arc::new(test_ctx(config::Config::default()).await);
|
||||||
|
let at = |path: &str| config::Feed { url: format!("{base}{path}"), ..feed() };
|
||||||
|
let none = db::HttpState::default();
|
||||||
|
let text = || ctx.db.strings_for_test("SELECT description FROM entries WHERE feed_id = 'aicn'");
|
||||||
|
scan_one(&ctx, "aicn", &at("/feed"), &none, false, None).await.unwrap();
|
||||||
|
assert_eq!(text().await, ["Somebody else's text."]);
|
||||||
|
// The feed fixed, the item already stored: it was left as it was (#141).
|
||||||
|
scan_one(&ctx, "aicn", &at("/fixed"), &none, false, None).await.unwrap();
|
||||||
|
assert_eq!(text().await, ["As it should read."]);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn nothing_from_a_feed_that_never_posted_is_an_error_and_from_one_that_has_is_not() {
|
||||||
|
let base = empty_or_feed_server().await;
|
||||||
|
let ctx = Arc::new(test_ctx(config::Config::default()).await);
|
||||||
|
let at = |path: &str| config::Feed { url: format!("{base}{path}"), ..feed() };
|
||||||
|
let none = db::HttpState::default();
|
||||||
|
// A lapsed domain behind a block page: never an item, and nothing (#119).
|
||||||
|
let err = scan_one(&ctx, "gone", &at("/empty"), &none, false, None).await.err().expect("an error");
|
||||||
|
assert!(format!("{err:#}").contains("answered with nothing"), "{err:#}");
|
||||||
|
// A feed that has posted answering nothing has nothing new, as the Antarctic Survey does.
|
||||||
|
assert!(matches!(scan_one(&ctx, "quiet", &at("/feed"), &none, false, None).await, Ok(Outcome::Feed(_))));
|
||||||
|
assert!(matches!(scan_one(&ctx, "quiet", &at("/empty"), &none, false, None).await, Ok(Outcome::Empty)));
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn a_derived_feed_is_not_scanned_once_its_opml_leaves_config() {
|
async fn a_derived_feed_is_not_scanned_once_its_opml_leaves_config() {
|
||||||
// davewiner: the OPML subscription left config.toml, but its 922 derived rows
|
// davewiner: the OPML subscription left config.toml, but its 922 derived rows
|
||||||
|
|||||||
26
src/web.rs
26
src/web.rs
@@ -315,11 +315,8 @@ async fn login(
|
|||||||
) -> Result<Response, ApiError> {
|
) -> Result<Response, ApiError> {
|
||||||
let name = body.name.trim().to_ascii_lowercase();
|
let name = body.name.trim().to_ascii_lowercase();
|
||||||
let user = state.ctx.db.user_by_name(&name).await?;
|
let user = state.ctx.db.user_by_name(&name).await?;
|
||||||
// The same answer either way: whether a name exists is not something to leak.
|
// The same answer either way, in the same time: whether a name exists is not something to leak.
|
||||||
let ok = user
|
let ok = crate::auth::check_password(body.password, user.as_ref().and_then(|u| u.pass_hash.clone())).await;
|
||||||
.as_ref()
|
|
||||||
.and_then(|u| u.pass_hash.as_deref())
|
|
||||||
.is_some_and(|h| crate::auth::verify_password(&body.password, h));
|
|
||||||
if !ok {
|
if !ok {
|
||||||
tracing::warn!(user = %name, "failed sign-in");
|
tracing::warn!(user = %name, "failed sign-in");
|
||||||
return Ok((StatusCode::UNAUTHORIZED, "wrong name or password").into_response());
|
return Ok((StatusCode::UNAUTHORIZED, "wrong name or password").into_response());
|
||||||
@@ -1027,9 +1024,17 @@ struct ListedItem {
|
|||||||
struct Listed {
|
struct Listed {
|
||||||
/// What the feed says it is (#130), sanitized as its items are.
|
/// What the feed says it is (#130), sanitized as its items are.
|
||||||
description: Option<String>,
|
description: Option<String>,
|
||||||
|
/// Why its last check failed, when it did (#142), in the words a feed's own page uses.
|
||||||
|
failing: Option<String>,
|
||||||
items: Vec<ListedItem>,
|
items: Vec<ListedItem>,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// A feed's last error as its Directory page says it: explained, or that the check failed, never
|
||||||
|
/// the error as it stands, which can name the feed's address, and the Directory names none.
|
||||||
|
fn failure_words(err: &str) -> String {
|
||||||
|
crate::feed::explain_failure(err).map_or_else(|| "Its last check failed.".to_owned(), |f| f.reason.to_owned())
|
||||||
|
}
|
||||||
|
|
||||||
/// A listed feed's description and latest twenty items. Only a feed the Directory lists, so a
|
/// A listed feed's description and latest twenty items. Only a feed the Directory lists, so a
|
||||||
/// guessed id reaches nothing private, as `subscribe_popular` checks.
|
/// guessed id reaches nothing private, as `subscribe_popular` checks.
|
||||||
async fn get_listed(
|
async fn get_listed(
|
||||||
@@ -1044,6 +1049,7 @@ async fn get_listed(
|
|||||||
let rows = state.ctx.db.entries_in(user.id, Some(&id), crate::db::Filter::parse("all"), None, 0, 20, &order).await?;
|
let rows = state.ctx.db.entries_in(user.id, Some(&id), crate::db::Filter::parse("all"), None, 0, 20, &order).await?;
|
||||||
let mut sanitizer = feed_sanitizer();
|
let mut sanitizer = feed_sanitizer();
|
||||||
let description = state.ctx.db.feed_description(&id).await?.map(|d| clean_description(&mut sanitizer, &d, None));
|
let description = state.ctx.db.feed_description(&id).await?.map(|d| clean_description(&mut sanitizer, &d, None));
|
||||||
|
let failing = state.ctx.db.feed_summary(&id).await?.last_error.as_deref().map(failure_words);
|
||||||
let items = rows
|
let items = rows
|
||||||
.into_iter()
|
.into_iter()
|
||||||
.map(|e| ListedItem {
|
.map(|e| ListedItem {
|
||||||
@@ -1055,7 +1061,7 @@ async fn get_listed(
|
|||||||
duration: e.duration,
|
duration: e.duration,
|
||||||
})
|
})
|
||||||
.collect();
|
.collect();
|
||||||
Ok(Json(Listed { description, items }))
|
Ok(Json(Listed { description, failing, items }))
|
||||||
}
|
}
|
||||||
|
|
||||||
fn sort_name(p: &PopularRow) -> String {
|
fn sort_name(p: &PopularRow) -> String {
|
||||||
@@ -1149,6 +1155,14 @@ impl IntoResponse for ApiError {
|
|||||||
mod tests {
|
mod tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn a_failing_feed_says_why_in_the_directory_without_its_address() {
|
||||||
|
assert_eq!(failure_words("the feed's address answered with nothing at all"), "This address answers with nothing; the site may be gone.");
|
||||||
|
// reqwest names the URL; the Directory names none.
|
||||||
|
let raw = "error sending request for url (http://daily-quests.com/comic/?feed=atom): connection closed";
|
||||||
|
assert_eq!(failure_words(raw), "Its last check failed.");
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn a_relative_image_resolves_against_the_post() {
|
fn a_relative_image_resolves_against_the_post() {
|
||||||
let mut b = feed_sanitizer();
|
let mut b = feed_sanitizer();
|
||||||
|
|||||||
54
tests/load/browse.js
Normal file
54
tests/load/browse.js
Normal file
@@ -0,0 +1,54 @@
|
|||||||
|
import http from 'k6/http';
|
||||||
|
import { check, sleep } from 'k6';
|
||||||
|
import { BASE, FEEDS, ITEMS, as, me, pick, feedId, p95 } from './lib.js';
|
||||||
|
|
||||||
|
// An evening's browsing: twenty-five people at once opening their feed list, All Subscriptions,
|
||||||
|
// a feed, a search, the Directory and a feed's page in it. Each answer's budget is well above what
|
||||||
|
// it takes now, so it fails on a query that has started asking once per feed -- the Directory
|
||||||
|
// asked the database three questions a feed until 0.10.0 -- not on a slow minute. Fifty measured
|
||||||
|
// the queue for what was SQLite's one connection (#136) more than any query.
|
||||||
|
export const options = {
|
||||||
|
scenarios: {
|
||||||
|
evening: {
|
||||||
|
executor: 'ramping-vus',
|
||||||
|
stages: [{ duration: '15s', target: 25 }, { duration: '45s', target: 25 }, { duration: '10s', target: 0 }],
|
||||||
|
},
|
||||||
|
},
|
||||||
|
thresholds: {
|
||||||
|
http_req_failed: ['rate==0'],
|
||||||
|
checks: ['rate==1'],
|
||||||
|
// About twice each one's p95 on 2026-10-05 on Tower, and no less than 50ms, where a few ms
|
||||||
|
// either way is noise. SQLite: 89, 17, 10, 22, 112, 126ms; Postgres: 54, 34, 24, 44, 71, 88ms.
|
||||||
|
'http_req_duration{name:feeds}': p95(200, 150),
|
||||||
|
'http_req_duration{name:all-entries}': p95(50, 100),
|
||||||
|
'http_req_duration{name:feed-entries}': p95(50, 100),
|
||||||
|
'http_req_duration{name:search}': p95(60, 100),
|
||||||
|
'http_req_duration{name:directory}': p95(250, 200),
|
||||||
|
'http_req_duration{name:listed}': p95(300, 200),
|
||||||
|
},
|
||||||
|
};
|
||||||
|
|
||||||
|
export default function () {
|
||||||
|
const u = me();
|
||||||
|
const feeds = http.get(`${BASE}/api/feeds`, as(u, 'feeds'));
|
||||||
|
check(feeds, { 'your thirty feeds': r => r.status === 200 && r.json().length === 30 });
|
||||||
|
|
||||||
|
const all = http.get(`${BASE}/api/entries?limit=50&filter=all&sort=published&dir=desc`, as(u, 'all-entries'));
|
||||||
|
check(all, { 'All Subscriptions, a page of it': r => r.status === 200 && r.json().entries.length === 50 });
|
||||||
|
|
||||||
|
const f = pick(feeds.json());
|
||||||
|
const one = http.get(`${BASE}/api/feeds/${f.id}/entries?limit=50&sort=title&dir=asc`, as(u, 'feed-entries'));
|
||||||
|
check(one, { 'a feed of yours, every item': r => r.status === 200 && r.json().entries.length === ITEMS });
|
||||||
|
|
||||||
|
// Every generated item mentions the weather, in its text, not its title.
|
||||||
|
const found = http.get(`${BASE}/api/entries?q=weather&limit=50`, as(u, 'search'));
|
||||||
|
check(found, { 'a search of everything you read': r => r.status === 200 && r.json().total === 30 * ITEMS });
|
||||||
|
|
||||||
|
const dir = http.get(`${BASE}/api/directory`, as(u, 'directory'));
|
||||||
|
check(dir, { 'the whole Directory': r => r.status === 200 && r.json().length === FEEDS });
|
||||||
|
|
||||||
|
const listed = http.get(`${BASE}/api/directory/${feedId(Math.floor(Math.random() * FEEDS))}`, as(u, 'listed'));
|
||||||
|
check(listed, { "a feed's page in the Directory": r => r.status === 200 && r.json().items.length === ITEMS });
|
||||||
|
|
||||||
|
sleep(1 + Math.random() * 2);
|
||||||
|
}
|
||||||
19
tests/load/lib.js
Normal file
19
tests/load/lib.js
Normal file
@@ -0,0 +1,19 @@
|
|||||||
|
// What the load tests share. Each virtual user is a listener of its own, signed in by name with
|
||||||
|
// the header the scratch daemon trusts from 127.0.0.1, as production trusts Cloudflare Access's.
|
||||||
|
export const BASE = __ENV.BASE;
|
||||||
|
export const USERS = Number(__ENV.USERS || 100);
|
||||||
|
export const FEEDS = Number(__ENV.FEEDS || 1500);
|
||||||
|
// Each generated feed has this many items; see run.js.
|
||||||
|
export const ITEMS = 20;
|
||||||
|
|
||||||
|
export const as = (name, tag) => ({
|
||||||
|
headers: { 'X-Load-User': name, 'Content-Type': 'application/json' },
|
||||||
|
tags: { name: tag },
|
||||||
|
});
|
||||||
|
/// The listener this virtual user is, one of the hundred run.js seeded, thirty feeds each.
|
||||||
|
export const me = () => `load-${(__VU - 1) % USERS}`;
|
||||||
|
export const pick = a => a[Math.floor(Math.random() * a.length)];
|
||||||
|
/// A p95 budget, in ms, for whichever database the daemon is on (run.js --postgres): SQLite reads
|
||||||
|
/// a file beside the daemon, Postgres answers over the network, and each is quicker at something.
|
||||||
|
export const p95 = (sqlite, postgres) => [`p(95)<${__ENV.DB === 'postgres' ? postgres : sqlite}`];
|
||||||
|
export const feedId = n => `gen-${String(n).padStart(4, '0')}`;
|
||||||
57
tests/load/listening.js
Normal file
57
tests/load/listening.js
Normal file
@@ -0,0 +1,57 @@
|
|||||||
|
import http from 'k6/http';
|
||||||
|
import { check, sleep } from 'k6';
|
||||||
|
import { BASE, as, me, pick, p95 } from './lib.js';
|
||||||
|
|
||||||
|
// Players saving where people are while scans write. A hundred listeners each save a position
|
||||||
|
// every second -- the player saves every ten, so this is a thousand people listening -- mark an
|
||||||
|
// item read now and then, and read their position back to see it is theirs and as they left it,
|
||||||
|
// while one of them forces a scan of all their feeds every five seconds. On SQLite every one of
|
||||||
|
// these is a write waiting its turn for the one writer.
|
||||||
|
export const options = {
|
||||||
|
scenarios: {
|
||||||
|
listeners: { executor: 'constant-vus', vus: 100, duration: '60s', exec: 'listen' },
|
||||||
|
scans: { executor: 'constant-vus', vus: 1, duration: '60s', exec: 'scan' },
|
||||||
|
},
|
||||||
|
thresholds: {
|
||||||
|
http_req_failed: ['rate==0'],
|
||||||
|
checks: ['rate==1'],
|
||||||
|
// About twice each one's p95 on 2026-10-05 on Tower. SQLite: 21, 277 and 134ms; Postgres: 9,
|
||||||
|
// 225 and 153ms. Marking read answers with the feed's row, counts and all, hence its budget.
|
||||||
|
'http_req_duration{name:position}': p95(100, 100),
|
||||||
|
'http_req_duration{name:flags}': p95(600, 500),
|
||||||
|
'http_req_duration{name:read-back}': p95(300, 400),
|
||||||
|
},
|
||||||
|
};
|
||||||
|
|
||||||
|
// Each virtual user's own episode, and how far into it it is.
|
||||||
|
let ep = null, secs = 0, n = 0;
|
||||||
|
|
||||||
|
export function listen() {
|
||||||
|
const u = me();
|
||||||
|
if (!ep) {
|
||||||
|
const r = http.get(`${BASE}/api/entries?limit=50`, as(u, 'pick'));
|
||||||
|
ep = pick(r.json().entries.filter(e => e.duration));
|
||||||
|
}
|
||||||
|
secs += 1;
|
||||||
|
const saved = http.post(`${BASE}/api/entries/${encodeURIComponent(ep.feed_id)}/${encodeURIComponent(ep.guid)}/position`,
|
||||||
|
JSON.stringify({ secs, duration: ep.duration }), as(u, 'position'));
|
||||||
|
check(saved, { 'position saved': r => r.status === 204 });
|
||||||
|
if (++n % 10 === 0) {
|
||||||
|
const flags = http.post(`${BASE}/api/entries/${encodeURIComponent(ep.feed_id)}/${encodeURIComponent(ep.guid)}/flags`,
|
||||||
|
JSON.stringify({ read: n % 20 === 0 }), as(u, 'flags'));
|
||||||
|
check(flags, { 'marked read or unread': r => r.status === 200 });
|
||||||
|
const back = http.get(`${BASE}/api/feeds/${encodeURIComponent(ep.feed_id)}/entries?limit=50`, as(u, 'read-back'));
|
||||||
|
const mine = back.status === 200 && back.json().entries.find(e => e.guid === ep.guid);
|
||||||
|
check(mine, {
|
||||||
|
'your position, as you left it': e => e && e.position === secs,
|
||||||
|
'read as you marked it': e => e && e.read === (n % 20 === 0),
|
||||||
|
});
|
||||||
|
}
|
||||||
|
sleep(1);
|
||||||
|
}
|
||||||
|
|
||||||
|
export function scan() {
|
||||||
|
const r = http.post(`${BASE}/api/fetch`, JSON.stringify({ force: true }), as('load-0', 'fetch'));
|
||||||
|
check(r, { 'scan queued': r => r.status === 202 || r.status === 200 });
|
||||||
|
sleep(5);
|
||||||
|
}
|
||||||
48
tests/load/media.js
Normal file
48
tests/load/media.js
Normal file
@@ -0,0 +1,48 @@
|
|||||||
|
import http from 'k6/http';
|
||||||
|
import { check } from 'k6';
|
||||||
|
import { BASE, as, pick } from './lib.js';
|
||||||
|
|
||||||
|
// Listeners seeking: fifty at once asking for ranges of the same few episodes, as a player does
|
||||||
|
// on every seek and every few seconds of playback. Every range is checked against what the feed
|
||||||
|
// served, byte by byte at a sample of offsets, so a wrong offset fails, not only a wrong status.
|
||||||
|
export const options = {
|
||||||
|
scenarios: { seeking: { executor: 'constant-vus', vus: 50, duration: '45s' } },
|
||||||
|
thresholds: {
|
||||||
|
http_req_failed: ['rate==0'],
|
||||||
|
checks: ['rate==1'],
|
||||||
|
// About four times its p95 on 2026-10-05, 48ms on SQLite and 45 on Postgres: a range is a
|
||||||
|
// file read, little to vary.
|
||||||
|
'http_req_duration{name:range}': ['p(95)<200'],
|
||||||
|
},
|
||||||
|
};
|
||||||
|
// The generated file: ID3's ten bytes, then a byte its offset gives (run.js, mediaByte).
|
||||||
|
const SIZE = 2 * 1024 * 1024;
|
||||||
|
const byte = k => (k * 31 + 7) & 255;
|
||||||
|
|
||||||
|
export function setup() {
|
||||||
|
const r = http.get(`${BASE}/api/entries?filter=downloaded&limit=50`, as('listener', 'setup'));
|
||||||
|
const files = r.json().entries.flatMap(e => e.enclosures).filter(x => x.path).map(x => x.id);
|
||||||
|
if (!files.length) throw new Error('the listener has no downloaded files to seek through');
|
||||||
|
return { files };
|
||||||
|
}
|
||||||
|
|
||||||
|
export default function ({ files }) {
|
||||||
|
const start = 10 + Math.floor(Math.random() * (SIZE - 11));
|
||||||
|
const end = Math.min(SIZE - 1, start + Math.floor(Math.random() * 65536));
|
||||||
|
const r = http.get(`${BASE}/media/${pick(files)}`, {
|
||||||
|
headers: { 'X-Load-User': 'listener', Range: `bytes=${start}-${end}` },
|
||||||
|
responseType: 'binary',
|
||||||
|
tags: { name: 'range' },
|
||||||
|
});
|
||||||
|
const body = r.status === 206 ? new Uint8Array(r.body) : new Uint8Array(0);
|
||||||
|
check(r, {
|
||||||
|
'part of the file': r => r.status === 206,
|
||||||
|
'the range asked for': r => r.headers['Content-Range'] === `bytes ${start}-${end}/${SIZE}`,
|
||||||
|
'its bytes': () => {
|
||||||
|
if (body.length !== end - start + 1) return false;
|
||||||
|
for (const at of [0, body.length - 1, ...Array.from({ length: 30 }, () => Math.floor(Math.random() * body.length))])
|
||||||
|
if (body[at] !== byte(start + at)) return false;
|
||||||
|
return true;
|
||||||
|
},
|
||||||
|
});
|
||||||
|
}
|
||||||
280
tests/load/run.js
Normal file
280
tests/load/run.js
Normal file
@@ -0,0 +1,280 @@
|
|||||||
|
// Load tests: k6 against a scratch daemon with a catalogue the size of production's, for what
|
||||||
|
// the browser tests cannot show -- many people at once, and what that does to latency, to the
|
||||||
|
// database and to the healthcheck. The browser suite drives one person against a handful of
|
||||||
|
// feeds, one request at a time.
|
||||||
|
//
|
||||||
|
// node tests/load/run.js [--postgres] [browse|listening|media|signin ...] default: all four
|
||||||
|
//
|
||||||
|
// It builds a release binary (debug Argon2 alone takes a second a sign-in, which would measure
|
||||||
|
// the build, not ipx), serves generated feeds from this process, starts a daemon on its own
|
||||||
|
// config and data under /tmp/ipx-load, wiped first, seeds listeners, then runs each k6 script.
|
||||||
|
// While each runs, `ipx status` -- the Docker healthcheck -- is run every second, and the run
|
||||||
|
// fails if any answer takes as long as the healthcheck's 5s timeout.
|
||||||
|
//
|
||||||
|
// The daemon is on SQLite unless --postgres, which puts it in IPX_TEST_DATABASE_URL, the
|
||||||
|
// database the Rust tests use on Postgres too (ipodderx_test on the server production's is on, in
|
||||||
|
// /src/.envrc), in a schema of its own, ipx_load, made afresh each run (#139). Postgres's numbers
|
||||||
|
// are what production would see.
|
||||||
|
// IPX_LOAD_FEEDS sets the catalogue's size (1500, near production's), IPX_LOAD_USERS the
|
||||||
|
// listeners (100), IPX_LOAD_LOG=1 shows the daemon's log.
|
||||||
|
const http = require('http');
|
||||||
|
const fs = require('fs');
|
||||||
|
const path = require('path');
|
||||||
|
const { spawn, spawnSync, execFile, execFileSync } = require('child_process');
|
||||||
|
|
||||||
|
const repo = path.resolve(__dirname, '../..');
|
||||||
|
const root = '/tmp/ipx-load';
|
||||||
|
const WEB = 8793, GEN = 8794;
|
||||||
|
const FEEDS = Number(process.env.IPX_LOAD_FEEDS || 1500);
|
||||||
|
const USERS = Number(process.env.IPX_LOAD_USERS || 100);
|
||||||
|
// The first few feeds carry real files, downloaded for the media test.
|
||||||
|
const MEDIA = 4;
|
||||||
|
const ITEMS = 20;
|
||||||
|
const BASE = `http://127.0.0.1:${WEB}`;
|
||||||
|
const bin = path.join(repo, 'target/release/ipx');
|
||||||
|
const SCRIPTS = ['browse', 'listening', 'media', 'signin'];
|
||||||
|
const POSTGRES = process.argv.includes('--postgres');
|
||||||
|
const PG_URL = process.env.IPX_TEST_DATABASE_URL || '';
|
||||||
|
const PG_SCHEMA = 'ipx_load';
|
||||||
|
const env = {
|
||||||
|
...process.env,
|
||||||
|
IPX_CONFIG: `${root}/config/config.toml`,
|
||||||
|
IPX_DATA_DIR: `${root}/data`,
|
||||||
|
IPX_LOG: 'ipx=info',
|
||||||
|
IPX_LOG_FORMAT: 'json',
|
||||||
|
// Nothing of the test reaches production's database or traces, or a paid API.
|
||||||
|
// In its own schema, as Db::memory puts each Rust test, and with notices off, as url_for does.
|
||||||
|
IPX_DATABASE_URL: POSTGRES
|
||||||
|
? `${PG_URL}${PG_URL.includes('?') ? '&' : '?'}options=-c%20search_path%3D${PG_SCHEMA}%20-c%20client_min_messages%3Dwarning`
|
||||||
|
: '',
|
||||||
|
OTEL_EXPORTER_OTLP_ENDPOINT: '',
|
||||||
|
TYPESAFE_KEY: '',
|
||||||
|
};
|
||||||
|
|
||||||
|
// ---- the feeds ------------------------------------------------------------------------------
|
||||||
|
|
||||||
|
// Apple's categories, some with a subcategory, so the Directory has tiles and pages to draw.
|
||||||
|
const CATS = [['Technology'], ['News', 'Tech News'], ['Comedy'], ['Leisure', 'Video Games'],
|
||||||
|
['Society & Culture', 'Documentary'], ['Sports', 'Soccer'], ['Arts', 'Books'],
|
||||||
|
['Science', 'Astronomy'], ['History'], ['True Crime'], ['Business', 'Investing'],
|
||||||
|
['Education', 'Self-Improvement'], ['Health & Fitness', 'Mental Health'], ['Music']];
|
||||||
|
const esc = s => s.replace(/&/g, '&').replace(/</g, '<').replace(/"/g, '"');
|
||||||
|
const now = Date.now();
|
||||||
|
// Every third feed is a blog: no files, as most of production's are.
|
||||||
|
const isBlog = n => n % 3 === 2 && n >= MEDIA;
|
||||||
|
function feedXml(n) {
|
||||||
|
const [cat, sub] = CATS[n % CATS.length];
|
||||||
|
const category = sub
|
||||||
|
? `<itunes:category text="${esc(cat)}"><itunes:category text="${esc(sub)}"/></itunes:category>`
|
||||||
|
: `<itunes:category text="${esc(cat)}"/>`;
|
||||||
|
const items = Array.from({ length: ITEMS }, (_, i) => `<item>
|
||||||
|
<title>Episode ${ITEMS - i} of Show ${n}</title><guid>gen-${n}-${i}</guid>
|
||||||
|
<pubDate>${new Date(now - (i * 3 + n % 3) * 86400e3).toUTCString()}</pubDate>
|
||||||
|
<description><p>Show ${n}, episode ${ITEMS - i}: an hour on the news, the weather and whatever came up.</p></description>
|
||||||
|
${isBlog(n) ? '' : `<itunes:duration>${1800 + i * 60}</itunes:duration>
|
||||||
|
<enclosure url="http://127.0.0.1:${GEN}/media/${n}/${i}.mp3" length="${MEDIA_BYTES.length}" type="audio/mpeg"/>`}
|
||||||
|
</item>`).join('');
|
||||||
|
return `<?xml version="1.0"?><rss version="2.0" xmlns:itunes="http://www.itunes.com/dtds/podcast-1.0.dtd"><channel>
|
||||||
|
<title>Generated Show ${n}</title><link>http://127.0.0.1:${GEN}/site/${n}</link>
|
||||||
|
<description>Generated show number ${n}, for the load tests.</description>
|
||||||
|
<itunes:image href="http://127.0.0.1:${GEN}/art/${n}.jpg"/>${category}${items}</channel></rss>`;
|
||||||
|
}
|
||||||
|
// A file whose every byte is known from its offset, so the media test can check that a range
|
||||||
|
// it asked for is the range it got, not merely that it got something that long. ID3 first, so
|
||||||
|
// ipx takes it for audio.
|
||||||
|
const mediaByte = k => (k * 31 + 7) & 255;
|
||||||
|
const MEDIA_BYTES = Buffer.from(Uint8Array.from({ length: 2 * 1024 * 1024 }, (_, k) => mediaByte(k)));
|
||||||
|
Buffer.from('ID3\x03\x00\x00\x00\x00\x00\x00', 'latin1').copy(MEDIA_BYTES);
|
||||||
|
const art = fs.readFileSync(path.join(repo, 'tests/ui/fixtures/art.jpg'));
|
||||||
|
|
||||||
|
function serveFeeds() {
|
||||||
|
return http.createServer((req, res) => {
|
||||||
|
const m = req.url.match(/^\/(feed|media|art)\/(\d+)/);
|
||||||
|
if (m?.[1] === 'feed') { res.writeHead(200, { 'content-type': 'application/rss+xml' }); return res.end(feedXml(Number(m[2]))); }
|
||||||
|
if (m?.[1] === 'media') { res.writeHead(200, { 'content-type': 'audio/mpeg', 'content-length': MEDIA_BYTES.length }); return res.end(MEDIA_BYTES); }
|
||||||
|
if (m?.[1] === 'art') { res.writeHead(200, { 'content-type': 'image/jpeg' }); return res.end(art); }
|
||||||
|
res.writeHead(404).end();
|
||||||
|
}).listen(GEN, '127.0.0.1');
|
||||||
|
}
|
||||||
|
const feedId = n => `gen-${String(n).padStart(4, '0')}`;
|
||||||
|
const feedUrl = n => `http://127.0.0.1:${GEN}/feed/${n}.xml`;
|
||||||
|
|
||||||
|
// ---- the daemon -----------------------------------------------------------------------------
|
||||||
|
|
||||||
|
/// The load tests' schema dropped and made again, as /tmp/ipx-load is wiped. Only in a database
|
||||||
|
/// made for tests: the URL is production's server and role, and one slip of a name would
|
||||||
|
/// otherwise point this at production's database.
|
||||||
|
function freshPostgres() {
|
||||||
|
const u = PG_URL && new URL(PG_URL);
|
||||||
|
if (!u || !/_test$/.test(u.pathname)) {
|
||||||
|
throw new Error('--postgres needs IPX_TEST_DATABASE_URL, a database whose name ends in _test (. /src/.envrc)');
|
||||||
|
}
|
||||||
|
psql(`DROP SCHEMA IF EXISTS ${PG_SCHEMA} CASCADE`, `CREATE SCHEMA ${PG_SCHEMA}`);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Statements run in the test database, given by its URL in psql's environment, not its
|
||||||
|
/// arguments, where every process on Tower could read the password.
|
||||||
|
function psql(...sql) {
|
||||||
|
const u = new URL(PG_URL);
|
||||||
|
const env = { ...process.env, PGHOST: u.hostname, PGPORT: u.port || '5432', PGUSER: decodeURIComponent(u.username),
|
||||||
|
PGPASSWORD: decodeURIComponent(u.password), PGDATABASE: u.pathname.slice(1), PGOPTIONS: `-c search_path=${PG_SCHEMA}` };
|
||||||
|
execFileSync('psql', ['-q', '-v', 'ON_ERROR_STOP=1', ...sql.flatMap(s => ['-c', s])], { env, stdio: ['ignore', 'ignore', 'inherit'] });
|
||||||
|
}
|
||||||
|
|
||||||
|
function writeConfig() {
|
||||||
|
fs.rmSync(root, { recursive: true, force: true });
|
||||||
|
for (const d of ['config', 'data', 'downloads']) fs.mkdirSync(path.join(root, d), { recursive: true });
|
||||||
|
// Read into the database's catalogue on the first start, as an existing config.toml is. Only
|
||||||
|
// the media feeds download: a forced scan of a listener's 30 feeds would otherwise fetch a
|
||||||
|
// 2 MB episode of each, every time.
|
||||||
|
const feeds = Array.from({ length: FEEDS }, (_, n) =>
|
||||||
|
`[feeds.${feedId(n)}]\nurl = "${feedUrl(n)}"\nauto_download = ${n < MEDIA}\n`).join('\n');
|
||||||
|
fs.writeFileSync(env.IPX_CONFIG, `
|
||||||
|
[general]
|
||||||
|
download_dir = "${root}/downloads"
|
||||||
|
socket = "${root}/ipx.sock"
|
||||||
|
schedule = "every 60m"
|
||||||
|
max_new_per_check = 1
|
||||||
|
|
||||||
|
[torrent]
|
||||||
|
enabled = false
|
||||||
|
|
||||||
|
[web]
|
||||||
|
enabled = true
|
||||||
|
bind = "127.0.0.1:${WEB}"
|
||||||
|
token = "loadtokenloadtokenloadtoken12345"
|
||||||
|
# Each k6 user signs in by name, as Cloudflare Access does in production: no password to hash
|
||||||
|
# for every one of a hundred listeners. Only the sign-in test uses a password.
|
||||||
|
trusted_header = "X-Load-User"
|
||||||
|
trusted_proxies = ["127.0.0.1"]
|
||||||
|
auto_create_users = true
|
||||||
|
|
||||||
|
${feeds}`);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Resolves with the first log line matching `test`, read from the daemon's JSON log, which it
|
||||||
|
/// writes to stderr.
|
||||||
|
function waitForLog(daemon, test, ms) {
|
||||||
|
return new Promise((resolve, reject) => {
|
||||||
|
let buf = '';
|
||||||
|
const timer = setTimeout(() => { daemon.stderr.off('data', on); reject(new Error(`no log line in ${ms / 1000}s for ${test}`)); }, ms);
|
||||||
|
const on = chunk => {
|
||||||
|
buf += chunk;
|
||||||
|
let i;
|
||||||
|
while ((i = buf.indexOf('\n')) >= 0) {
|
||||||
|
const line = buf.slice(0, i); buf = buf.slice(i + 1);
|
||||||
|
let ev; try { ev = JSON.parse(line); } catch { continue; }
|
||||||
|
if (test(ev)) { clearTimeout(timer); daemon.stderr.off('data', on); return resolve(ev); }
|
||||||
|
}
|
||||||
|
};
|
||||||
|
daemon.stderr.on('data', on);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
const as = name => ({ 'X-Load-User': name, 'Content-Type': 'application/json' });
|
||||||
|
async function api(name, url, opts = {}) {
|
||||||
|
const r = await fetch(BASE + url, { ...opts, headers: { ...as(name), ...opts.headers } });
|
||||||
|
if (!r.ok) throw new Error(`${opts.method || 'GET'} ${url} as ${name}: ${r.status} ${await r.text()}`);
|
||||||
|
return r.status === 204 ? null : r.json().catch(() => null);
|
||||||
|
}
|
||||||
|
const opml = ns => `<?xml version="1.0"?><opml version="2.0"><head><title>load</title></head><body>${
|
||||||
|
ns.map(n => `<outline type="rss" text="${feedId(n)}" xmlUrl="${feedUrl(n)}"/>`).join('')}</body></opml>`;
|
||||||
|
|
||||||
|
async function seed() {
|
||||||
|
// The first account made is the admin.
|
||||||
|
await api('admin', '/api/me');
|
||||||
|
// Every listener subscribes to 30 feeds, overlapping as people's do; one OPML each, one scan.
|
||||||
|
for (let u = 0; u < USERS; u++) {
|
||||||
|
const mine = Array.from({ length: 30 }, (_, k) => (u * 7 + k * 41) % FEEDS);
|
||||||
|
await api(`load-${u}`, '/api/opml', { method: 'POST', body: JSON.stringify({ xml: opml(mine) }) });
|
||||||
|
}
|
||||||
|
// The media test's listener takes the feeds with files. They may be downloaded already: the
|
||||||
|
// first start subscribes the admin to the whole catalogue, so the first scan fetched them.
|
||||||
|
for (let n = 0; n < MEDIA; n++) await api('listener', `/api/popular/${feedId(n)}`, { method: 'POST' });
|
||||||
|
for (let t = Date.now(); ; await new Promise(r => setTimeout(r, 500))) {
|
||||||
|
const got = await api('listener', '/api/entries?filter=downloaded&limit=50');
|
||||||
|
if (got.entries.flatMap(e => e.enclosures).filter(x => x.path).length >= MEDIA) break;
|
||||||
|
if (Date.now() - t > 120_000) throw new Error(`the listener's ${MEDIA} files did not download in 120s`);
|
||||||
|
}
|
||||||
|
// A password, for the sign-in test; the CLI hashes it as `ipx user add` always does.
|
||||||
|
execFileSync(bin, ['user', 'add', 'piper'], { input: 'piperpassword', env });
|
||||||
|
}
|
||||||
|
|
||||||
|
// ---- the run --------------------------------------------------------------------------------
|
||||||
|
|
||||||
|
/// `ipx status` once a second until stopped: the slowest answer, and any that failed. Not
|
||||||
|
/// spawnSync: blocked in it, this process stops draining the daemon's log, the pipe fills, and
|
||||||
|
/// the daemon stalls writing its next line, which is the test measuring itself.
|
||||||
|
function watchHealth() {
|
||||||
|
const seen = { slowest: 0, failed: 0 };
|
||||||
|
let on = true;
|
||||||
|
(async () => {
|
||||||
|
while (on) {
|
||||||
|
const t = Date.now();
|
||||||
|
const ok = await new Promise(res => execFile(bin, ['status'], { env, timeout: 5000 }, err => res(!err)));
|
||||||
|
seen.slowest = Math.max(seen.slowest, Date.now() - t);
|
||||||
|
if (!ok) seen.failed++;
|
||||||
|
await new Promise(res => setTimeout(res, 1000));
|
||||||
|
}
|
||||||
|
})();
|
||||||
|
return () => { on = false; return seen; };
|
||||||
|
}
|
||||||
|
|
||||||
|
function k6(script) {
|
||||||
|
return new Promise(resolve => {
|
||||||
|
const run = spawn('k6', ['run', '--quiet', '-e', `BASE=${BASE}`, '-e', `USERS=${USERS}`, '-e', `FEEDS=${FEEDS}`,
|
||||||
|
'-e', `DB=${POSTGRES ? 'postgres' : 'sqlite'}`,
|
||||||
|
path.join(__dirname, `${script}.js`)], { stdio: 'inherit' });
|
||||||
|
run.on('exit', code => resolve(code));
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
(async () => {
|
||||||
|
const only = process.argv.slice(2).filter(a => a !== '--postgres');
|
||||||
|
for (const s of only) if (!SCRIPTS.includes(s)) { console.error(`no load test called ${s}: ${SCRIPTS.join(', ')}`); process.exit(2); }
|
||||||
|
if (spawnSync('k6', ['version']).error) { console.error('k6 is not installed: see install.sh'); process.exit(2); }
|
||||||
|
console.log('building the release binary...');
|
||||||
|
execFileSync('cargo', ['build', '--release', '-q'], { cwd: repo, stdio: 'inherit' });
|
||||||
|
|
||||||
|
const feeds = serveFeeds();
|
||||||
|
writeConfig();
|
||||||
|
if (POSTGRES) freshPostgres();
|
||||||
|
console.log(`on ${POSTGRES ? 'Postgres' : 'SQLite'}`);
|
||||||
|
// Its log kept out of this run's output, which is k6's; IPX_LOAD_LOG=1 shows it.
|
||||||
|
const daemon = spawn(bin, ['daemon'], { env, stdio: ['ignore', process.env.IPX_LOAD_LOG ? 'inherit' : 'ignore', 'pipe'] });
|
||||||
|
daemon.stderr.setEncoding('utf8');
|
||||||
|
if (process.env.IPX_LOAD_LOG) daemon.stderr.on('data', d => process.stderr.write(d));
|
||||||
|
// Stopped by hand, the daemon goes too, by its own PID, or it holds the ports for the next run.
|
||||||
|
for (const sig of ['SIGINT', 'SIGTERM']) process.on(sig, () => { daemon.kill('SIGTERM'); process.exit(130); });
|
||||||
|
let failed = 0;
|
||||||
|
try {
|
||||||
|
const t = Date.now();
|
||||||
|
const first = await waitForLog(daemon, ev => ev.ev === 'scan_done' && ev.feeds > 0, 600_000);
|
||||||
|
console.log(`first scan: ${first.feeds} feeds in ${((Date.now() - t) / 1000).toFixed(1)}s`);
|
||||||
|
await seed();
|
||||||
|
// The planner's statistics, as production's long-standing tables have them. A schema filled
|
||||||
|
// seconds ago has none until autovacuum gets to it, which it did partway through a test, and
|
||||||
|
// the search's p95 swung from 40ms to 166ms between runs with the plan it changed.
|
||||||
|
if (POSTGRES) psql('ANALYZE');
|
||||||
|
// The log is drained from here on, or the pipe fills and the daemon blocks writing to it.
|
||||||
|
daemon.stderr.resume();
|
||||||
|
for (const script of only.length ? only : SCRIPTS) {
|
||||||
|
console.log(`\n=== ${script}`);
|
||||||
|
const stop = watchHealth();
|
||||||
|
const code = await k6(script);
|
||||||
|
const health = stop();
|
||||||
|
console.log(`ipx status: slowest ${health.slowest}ms, ${health.failed} failed`);
|
||||||
|
if (code !== 0) { failed++; console.log(`${script}: thresholds failed (k6 exit ${code})`); }
|
||||||
|
if (health.failed || health.slowest >= 5000) { failed++; console.log(`${script}: the healthcheck would have failed`); }
|
||||||
|
}
|
||||||
|
} catch (e) {
|
||||||
|
console.error(e.message);
|
||||||
|
failed++;
|
||||||
|
} finally {
|
||||||
|
// Its own PID, never a pkill: Tower sees production's ipx too (#38).
|
||||||
|
daemon.kill('SIGTERM');
|
||||||
|
feeds.close();
|
||||||
|
}
|
||||||
|
console.log(failed ? `\n${failed} failure(s)` : '\nall load tests passed');
|
||||||
|
process.exit(failed ? 1 : 0);
|
||||||
|
})();
|
||||||
44
tests/load/signin.js
Normal file
44
tests/load/signin.js
Normal file
@@ -0,0 +1,44 @@
|
|||||||
|
import http from 'k6/http';
|
||||||
|
import { check, sleep } from 'k6';
|
||||||
|
import { Trend } from 'k6/metrics';
|
||||||
|
import { BASE, as, me } from './lib.js';
|
||||||
|
|
||||||
|
// A flood of sign-in attempts with wrong passwords while everyone else goes on using the site.
|
||||||
|
// Two things only many requests at once show: whether checking passwords, tens of milliseconds
|
||||||
|
// of CPU each, starves the rest of the server, and whether a wrong password for a name that is
|
||||||
|
// an account takes longer to refuse than one for a name that is not -- the answers read the
|
||||||
|
// same, and the time would tell anyone which names are accounts here.
|
||||||
|
const gap = new Trend('signin_name_gap', true);
|
||||||
|
export const options = {
|
||||||
|
scenarios: {
|
||||||
|
attempts: { executor: 'constant-vus', vus: 40, duration: '40s', exec: 'attempt' },
|
||||||
|
everyone: { executor: 'constant-vus', vus: 5, duration: '40s', exec: 'browse' },
|
||||||
|
},
|
||||||
|
thresholds: {
|
||||||
|
http_req_failed: ['rate==0'],
|
||||||
|
checks: ['rate==1'],
|
||||||
|
// 57ms on 2026-10-05 on SQLite, 68 on Postgres; 4.3s while password checks ran on the async
|
||||||
|
// workers (#137).
|
||||||
|
'http_req_duration{name:meanwhile}': ['p(95)<300'],
|
||||||
|
// Known name against unknown, one straight after the other: no gap but noise. 0.2ms on
|
||||||
|
// 2026-10-05; 31ms while an unknown name was refused without a check (#138).
|
||||||
|
signin_name_gap: ['med<10'],
|
||||||
|
},
|
||||||
|
};
|
||||||
|
const refused = { headers: { 'Content-Type': 'application/json' }, responseCallback: http.expectedStatuses(401) };
|
||||||
|
|
||||||
|
export function attempt() {
|
||||||
|
const known = http.post(`${BASE}/api/login`, JSON.stringify({ name: 'piper', password: 'not-the-password' }),
|
||||||
|
{ ...refused, tags: { name: 'signin-known' } });
|
||||||
|
const unknown = http.post(`${BASE}/api/login`, JSON.stringify({ name: `nobody-${__VU}-${__ITER}`, password: 'not-the-password' }),
|
||||||
|
{ ...refused, tags: { name: 'signin-unknown' } });
|
||||||
|
check(known, { 'a wrong password refused': r => r.status === 401 });
|
||||||
|
check(unknown, { 'an unknown name refused': r => r.status === 401 });
|
||||||
|
gap.add(known.timings.duration - unknown.timings.duration);
|
||||||
|
}
|
||||||
|
|
||||||
|
export function browse() {
|
||||||
|
const r = http.get(`${BASE}/api/feeds`, as(me(), 'meanwhile'));
|
||||||
|
check(r, { 'the site still answers': r => r.status === 200 });
|
||||||
|
sleep(0.5);
|
||||||
|
}
|
||||||
@@ -263,6 +263,15 @@ test('Currently Listening, its own place below the Directory, resumes an episode
|
|||||||
await expect(row).toBeVisible({ timeout: 20_000 });
|
await expect(row).toBeVisible({ timeout: 20_000 });
|
||||||
await expect(row).toContainText('14:18 left');
|
await expect(row).toContainText('14:18 left');
|
||||||
|
|
||||||
|
// The search box looks through it, as it does a feed's items (#127).
|
||||||
|
await page.locator('#epSearch').fill('second');
|
||||||
|
await expect(row).toBeVisible();
|
||||||
|
await page.locator('#epSearch').fill('nothing-is-called-this');
|
||||||
|
await expect(row).toHaveCount(0);
|
||||||
|
await expect(page.locator('#listening')).toContainText('Nothing you are listening to mentions');
|
||||||
|
await page.locator('#epSearch').fill('');
|
||||||
|
await expect(row).toBeVisible();
|
||||||
|
|
||||||
// Removing it forgets where you got to, so it is still gone on the next visit.
|
// Removing it forgets where you got to, so it is still gone on the next visit.
|
||||||
await row.locator('[data-a=remove]').click();
|
await row.locator('[data-a=remove]').click();
|
||||||
await expect(row).toHaveCount(0);
|
await expect(row).toHaveCount(0);
|
||||||
@@ -966,6 +975,46 @@ test('the Directory lists what everyone here reads, the most subscribed first, b
|
|||||||
await ctx.close();
|
await ctx.close();
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test('the feed list keeps the icons it has drawn when a feed\'s new row comes in (#134)', async ({ page }) => {
|
||||||
|
await expect(page.locator('#feedlist .feed').first()).toBeVisible({ timeout: 20_000 });
|
||||||
|
const kept = await page.evaluate(async () => {
|
||||||
|
const frames = () => new Promise(r => requestAnimationFrame(() => requestAnimationFrame(r)));
|
||||||
|
// Artwork that loads, on a feed of its own row: the fixtures' fails on purpose, for initials.
|
||||||
|
const row = { ...S.feeds.find(f => !f.group && !S.feeds.some(c => c.group === f.id)), image: '/favicon.png' };
|
||||||
|
patchFeed(row); await frames();
|
||||||
|
const rows = [...document.querySelectorAll('#feedlist .feed')];
|
||||||
|
const before = [...document.querySelectorAll('#feedlist .feed img')];
|
||||||
|
// The same feed's new row again, as a check of every feed sends one for each feed it reads.
|
||||||
|
patchFeed({ ...row }); await frames();
|
||||||
|
const after = [...document.querySelectorAll('#feedlist .feed img')];
|
||||||
|
return { drawn: before.length > 0, redrawn: !rows[0].isConnected, same: after.length === before.length && after.every((img, i) => img === before[i]) };
|
||||||
|
});
|
||||||
|
// The rows are new, the images in them the ones already drawn.
|
||||||
|
expect(kept).toEqual({ drawn: true, redrawn: true, same: true });
|
||||||
|
});
|
||||||
|
|
||||||
|
test('ipx add and rm while the daemon runs reach it, and its next change keeps them (#120)', async ({ page }) => {
|
||||||
|
const { execFileSync } = require('child_process');
|
||||||
|
const setup = require('./global-setup');
|
||||||
|
const env = { ...process.env, IPX_CONFIG: `${setup.root}/config/config.toml`, IPX_DATA_DIR: `${setup.root}/data` };
|
||||||
|
const ipx = (...args) => execFileSync('./target/debug/ipx', args, { env, encoding: 'utf8' });
|
||||||
|
const listed = async () => (await page.evaluate(() => api('/api/directory'))).map(p => p.id);
|
||||||
|
// A change of the daemon's own: it writes the whole catalogue from the copy it holds, which
|
||||||
|
// wrote a removed feed back and dropped an added one when that copy was old.
|
||||||
|
const daemonWrites = category => page.evaluate(c =>
|
||||||
|
api('/api/feeds/picture-blog', { method: 'PATCH', body: JSON.stringify({ category: c }) }), category);
|
||||||
|
|
||||||
|
const id = ipx('add', '--list', 'http://127.0.0.1:8792/cli.xml').match(/^added (\S+?),? /m)[1];
|
||||||
|
await expect.poll(listed).toContain(id);
|
||||||
|
await daemonWrites('Comedy');
|
||||||
|
expect(await listed()).toContain(id);
|
||||||
|
|
||||||
|
ipx('rm', id);
|
||||||
|
await expect.poll(listed).not.toContain(id);
|
||||||
|
await daemonWrites(null);
|
||||||
|
expect(await listed()).not.toContain(id);
|
||||||
|
});
|
||||||
|
|
||||||
test('adding a feed scans it straight away', async ({ page }) => {
|
test('adding a feed scans it straight away', async ({ page }) => {
|
||||||
await page.locator('#addFeed').click();
|
await page.locator('#addFeed').click();
|
||||||
await page.locator('#nurl').fill('http://127.0.0.1:8792/fresh.xml');
|
await page.locator('#nurl').fill('http://127.0.0.1:8792/fresh.xml');
|
||||||
|
|||||||
5
tests/ui/fixtures/cli.xml
Normal file
5
tests/ui/fixtures/cli.xml
Normal file
@@ -0,0 +1,5 @@
|
|||||||
|
<?xml version="1.0"?>
|
||||||
|
<rss version="2.0"><channel><title>Command Line Show</title><link>http://127.0.0.1:8792/</link>
|
||||||
|
<description>Added with ipx add while the daemon runs.</description>
|
||||||
|
<item><title>From the Command Line</title><guid>cli-1</guid><description>x</description></item>
|
||||||
|
</channel></rss>
|
||||||
@@ -717,6 +717,7 @@ body.scan-this .fhead [data-a=scan] .i,body.scan-any .fhead [data-a=scanall] .i,
|
|||||||
.show .meta{min-width:0;flex:1;display:flex;flex-direction:column;align-items:flex-start;gap:5px}
|
.show .meta{min-width:0;flex:1;display:flex;flex-direction:column;align-items:flex-start;gap:5px}
|
||||||
.show h2{margin:0 0 2px;font-size:30px;line-height:1.15;font-weight:700;letter-spacing:-.02em;overflow-wrap:anywhere}
|
.show h2{margin:0 0 2px;font-size:30px;line-height:1.15;font-weight:700;letter-spacing:-.02em;overflow-wrap:anywhere}
|
||||||
.show .sub{margin:0;color:var(--dim);font-size:13px;font-variant-numeric:tabular-nums}
|
.show .sub{margin:0;color:var(--dim);font-size:13px;font-variant-numeric:tabular-nums}
|
||||||
|
.show .dfail{color:var(--bad)}
|
||||||
.show .acts{margin-top:8px}
|
.show .acts{margin-top:8px}
|
||||||
.dlink{padding:0;font-size:13px;color:var(--accent);text-align:left}
|
.dlink{padding:0;font-size:13px;color:var(--accent);text-align:left}
|
||||||
.dlink:hover{text-decoration:underline}
|
.dlink:hover{text-decoration:underline}
|
||||||
|
|||||||
@@ -138,6 +138,7 @@ function drawDirectory(box=$('#directory'), focus?: string){
|
|||||||
// Its category's page is one click away from here.
|
// Its category's page is one click away from here.
|
||||||
show.category?`<button type="button" class="dlink" data-go="${esc(show.category)}">${esc(show.subcategory||show.category)}</button>`:''}
|
show.category?`<button type="button" class="dlink" data-go="${esc(show.category)}">${esc(show.subcategory||show.category)}</button>`:''}
|
||||||
<p class="sub">${n?`${n} ${n===1?'person here subscribes':'people here subscribe'}`:'Nobody here subscribes yet'}</p>
|
<p class="sub">${n?`${n} ${n===1?'person here subscribes':'people here subscribe'}`:'Nobody here subscribes yet'}</p>
|
||||||
|
<p class="sub dfail" id="dfail" hidden></p>
|
||||||
<div class="acts"><button class="btn ico primary" data-a="sub" title="Subscribe" aria-label="Subscribe">${ICON.plus}</button></div></div></div>
|
<div class="acts"><button class="btn ico primary" data-a="sub" title="Subscribe" aria-label="Subscribe">${ICON.plus}</button></div></div></div>
|
||||||
<div class="about" id="dabout" hidden><p></p><button type="button" class="dlink" aria-expanded="false" hidden>More</button></div>`
|
<div class="about" id="dabout" hidden><p></p><button type="button" class="dlink" aria-expanded="false" hidden>More</button></div>`
|
||||||
+sec(show.podcast?'Latest episodes':'Latest posts','<div class="latest" id="dshow"><p class="hint">Loading…</p></div>');
|
+sec(show.podcast?'Latest episodes':'Latest posts','<div class="latest" id="dshow"><p class="hint">Loading…</p></div>');
|
||||||
@@ -217,6 +218,8 @@ async function fillShow(box,p){
|
|||||||
more.textContent=open?'Less':'More'; more.setAttribute('aria-expanded',String(open));
|
more.textContent=open?'Less':'More'; more.setAttribute('aria-expanded',String(open));
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
// A broken feed says so, or it reads as a quiet one (#142).
|
||||||
|
if(got.failing){ const f=$('#dfail',box); f.textContent=got.failing; f.hidden=false; }
|
||||||
const items=got.items;
|
const items=got.items;
|
||||||
list.innerHTML=items.length?items.map(e=>{
|
list.innerHTML=items.length?items.map(e=>{
|
||||||
const name=entryName(e), text=name.derived?'':plainText(e.description).slice(0,500);
|
const name=entryName(e), text=name.derived?'':plainText(e.description).slice(0,500);
|
||||||
@@ -224,7 +227,8 @@ async function fillShow(box,p){
|
|||||||
return `<article class="lep"><b>${e.link?`<a href="${esc(e.link)}" target="_blank" rel="noopener">${esc(name.text)}</a>`:esc(name.text)}</b>`
|
return `<article class="lep"><b>${e.link?`<a href="${esc(e.link)}" target="_blank" rel="noopener">${esc(name.text)}</a>`:esc(name.text)}</b>`
|
||||||
+(text?`<p>${esc(text)}</p>`:'')
|
+(text?`<p>${esc(text)}</p>`:'')
|
||||||
+(when.length?`<small>${when.map(s=>`<span>${s}</span>`).join('<span class="dot"></span>')}</small>`:'')+'</article>';
|
+(when.length?`<small>${when.map(s=>`<span>${s}</span>`).join('<span class="dot"></span>')}</small>`:'')+'</article>';
|
||||||
}).join(''):`<p class="hint">Nothing read from it yet. A feed nobody here subscribes to is checked once a day.</p>`;
|
}).join(''):`<p class="hint">Nothing read from it yet.${p.subscribers||got.failing?''
|
||||||
|
:' A feed nobody here subscribes to is checked once a day.'}</p>`;
|
||||||
}
|
}
|
||||||
|
|
||||||
/// A category as a tile: its name, how many feeds, and the covers of its three most subscribed
|
/// A category as a tile: its name, how many feeds, and the covers of its three most subscribed
|
||||||
@@ -259,7 +263,8 @@ async function renderListed(v){
|
|||||||
</div>
|
</div>
|
||||||
<div class="childlist" id="listening"><p class="hint">Loading…</p></div>`;
|
<div class="childlist" id="listening"><p class="hint">Loading…</p></div>`;
|
||||||
$('#count').textContent=v.title;
|
$('#count').textContent=v.title;
|
||||||
const n=await renderListening(v.url,$('#listening',box));
|
// The search box looks through these too, title and text, as it does a feed's items (#127).
|
||||||
|
const n=await renderListening(S.q?`${v.url}&q=${encodeURIComponent(S.q)}`:v.url,$('#listening',box));
|
||||||
if(VIEWS[S.feed]===v) $('#count').textContent=`${v.title}: ${plural(n,'episode')}`;
|
if(VIEWS[S.feed]===v) $('#count').textContent=`${v.title}: ${plural(n,'episode')}`;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -269,7 +274,8 @@ async function renderListed(v){
|
|||||||
async function renderListening(url,box){
|
async function renderListening(url,box){
|
||||||
let rows=[];
|
let rows=[];
|
||||||
try{ rows=(await api(url)).entries||[]; }catch{}
|
try{ rows=(await api(url)).entries||[]; }catch{}
|
||||||
box.innerHTML=rows.length?'':'<p class="hint">Nothing in progress. Episodes you start and do not finish show up here.</p>';
|
box.innerHTML=rows.length?'':S.q?`<p class="hint">Nothing you are listening to mentions “${esc(S.q)}”.</p>`
|
||||||
|
:'<p class="hint">Nothing in progress. Episodes you start and do not finish show up here.</p>';
|
||||||
for(const e of rows){
|
for(const e of rows){
|
||||||
// Carries its episode, for paintListenRow to repaint as the player moves.
|
// Carries its episode, for paintListenRow to repaint as the player moves.
|
||||||
const el: HTMLDivElement & {entry?: any}=document.createElement('div');
|
const el: HTMLDivElement & {entry?: any}=document.createElement('div');
|
||||||
@@ -329,9 +335,12 @@ $('#tbRead').onclick=()=>{ const e=cur(); if(e) epAction('read',e,null); };
|
|||||||
$('#tbFlag').onclick=()=>{ const e=cur(); if(e) epAction('flag',e,null); };
|
$('#tbFlag').onclick=()=>{ const e=cur(); if(e) epAction('flag',e,null); };
|
||||||
let searchT;
|
let searchT;
|
||||||
$('#epSearch').oninput=ev=>{ clearTimeout(searchT);
|
$('#epSearch').oninput=ev=>{ clearTimeout(searchT);
|
||||||
// The Directory looks through what it already has; the rest ask the server (#126).
|
// The Directory looks through what it already has; Currently Listening and the rest ask the
|
||||||
|
// server (#126, #127).
|
||||||
searchT=setTimeout(()=>{ S.q=ev.target.value; S.offset=0;
|
searchT=setTimeout(()=>{ S.q=ev.target.value; S.offset=0;
|
||||||
if(S.feed===':directory'){ dirShow=null; drawDirectory(); } else loadEntries(); },250); };
|
if(S.feed===':directory'){ dirShow=null; drawDirectory(); }
|
||||||
|
else if(S.feed===':listening') renderListed(VIEWS[':listening']);
|
||||||
|
else loadEntries(); },250); };
|
||||||
// Crossing the phone breakpoint moves the files between their pane and the text.
|
// Crossing the phone breakpoint moves the files between their pane and the text.
|
||||||
window.matchMedia?.('(max-width:820px)')?.addEventListener?.('change',()=>{ const e=cur(); if(e) showDetail(e); });
|
window.matchMedia?.('(max-width:820px)')?.addEventListener?.('change',()=>{ const e=cur(); if(e) showDetail(e); });
|
||||||
|
|
||||||
|
|||||||
@@ -66,6 +66,11 @@ function renderFeeds(){
|
|||||||
const row=back&&list.querySelector(`[data-id="${CSS.escape(back[0])}"]`);
|
const row=back&&list.querySelector(`[data-id="${CSS.escape(back[0])}"]`);
|
||||||
if(row) (back[1]&&$('.chev',row)||row).focus();
|
if(row) (back[1]&&$('.chev',row)||row).focus();
|
||||||
};
|
};
|
||||||
|
// The images already drawn, to go into the rows that replace theirs. Made anew, every icon in
|
||||||
|
// the list went blank and loaded again, once a frame through a check of every feed, as each
|
||||||
|
// feed's new row came in (#134). Matched by their markup, so a changed one is made anew.
|
||||||
|
const drawn=new Map<string,HTMLImageElement[]>();
|
||||||
|
for(const img of $$('img',list)){ const k=img.outerHTML; if(!drawn.has(k)) drawn.set(k,[]); drawn.get(k).push(img); }
|
||||||
list.innerHTML='';
|
list.innerHTML='';
|
||||||
const unreadAll=S.feeds.reduce((n,f)=>n+(f.unread||0),0);
|
const unreadAll=S.feeds.reduce((n,f)=>n+(f.unread||0),0);
|
||||||
const places=document.createElement('div');
|
const places=document.createElement('div');
|
||||||
@@ -139,6 +144,7 @@ function renderFeeds(){
|
|||||||
`<span class="badge${unread?'':' zero'}" title="${unread} unread">${unread>999?'999+':unread}</span>`;
|
`<span class="badge${unread?'':' zero'}" title="${unread} unread">${unread>999?'999+':unread}</span>`;
|
||||||
el.onclick=()=>{ selectFeed(f.id); nav(false); };
|
el.onclick=()=>{ selectFeed(f.id); nav(false); };
|
||||||
if(kids) $('.chev',el).onclick=ev=>{ ev.stopPropagation(); toggleGroup(f.id); };
|
if(kids) $('.chev',el).onclick=ev=>{ ev.stopPropagation(); toggleGroup(f.id); };
|
||||||
|
for(const img of $$('img',el)){ const was=drawn.get(img.outerHTML)?.shift(); if(was) img.replaceWith(was); }
|
||||||
list.appendChild(el);
|
list.appendChild(el);
|
||||||
}
|
}
|
||||||
paintScanning();
|
paintScanning();
|
||||||
|
|||||||
@@ -179,7 +179,7 @@ function keysModal(){
|
|||||||
[g('a'),'All Subscriptions'],[g('d'),'Directory'],
|
[g('a'),'All Subscriptions'],[g('d'),'Directory'],
|
||||||
[g('l'),'Currently Listening'],[g('s'),'Settings'],
|
[g('l'),'Currently Listening'],[g('s'),'Settings'],
|
||||||
[k('Shift')+' '+k('J'),'Next feed'],[k('Shift')+' '+k('K'),'Previous feed'],
|
[k('Shift')+' '+k('J'),'Next feed'],[k('Shift')+' '+k('K'),'Previous feed'],
|
||||||
[k('/'),'Search items, or the Directory'],[k('r'),'Refresh'],[k('['),'Show or hide the feed list'],
|
[k('/'),'Search items, Currently Listening or the Directory'],[k('r'),'Refresh'],[k('['),'Show or hide the feed list'],
|
||||||
['Items'],
|
['Items'],
|
||||||
[k('j')+' or '+k('n'),'Next item'],[k('k')+' or '+k('p'),'Previous item'],
|
[k('j')+' or '+k('n'),'Next item'],[k('k')+' or '+k('p'),'Previous item'],
|
||||||
[k('Shift')+' '+k('A'),'Mark all read'],
|
[k('Shift')+' '+k('A'),'Mark all read'],
|
||||||
|
|||||||
Reference in New Issue
Block a user