Compare commits
12 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 50c3b9db64 | |||
| 39ae43c80e | |||
| c9964778cb | |||
| a553d050ca | |||
| 9c2f3eaf1b | |||
| c63f70ad0c | |||
| bbdcbe213f | |||
| 8c4a5f396b | |||
| b43ba89778 | |||
| f1d360420c | |||
| 697e907c86 | |||
| 894dbbe31d |
19
CHANGELOG.md
19
CHANGELOG.md
@@ -7,6 +7,22 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
|||||||
|
|
||||||
## [Unreleased]
|
## [Unreleased]
|
||||||
|
|
||||||
|
## [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 +751,8 @@ 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.1...main
|
||||||
|
[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-dev"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"ammonia",
|
"ammonia",
|
||||||
"anyhow",
|
"anyhow",
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "ipx"
|
name = "ipx"
|
||||||
version = "0.10.0"
|
version = "0.10.2-dev"
|
||||||
edition = "2024"
|
edition = "2024"
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
|
|||||||
@@ -146,6 +146,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
|
||||||
|
|||||||
@@ -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();
|
||||||
|
|||||||
49
src/db.rs
49
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) })
|
||||||
}
|
}
|
||||||
@@ -2106,13 +2123,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();
|
||||||
}
|
}
|
||||||
let row = db.orm.query_one_raw(Statement::from_string(sea_orm::DbBackend::Sqlite, "PRAGMA synchronous")).await.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)");
|
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)");
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
|
|||||||
@@ -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));
|
||||||
|
|
||||||
|
|||||||
94
src/main.rs
94
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);
|
||||||
}
|
}
|
||||||
@@ -2248,6 +2282,60 @@ 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, anything else with a feed of one item.
|
||||||
|
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 empty = String::from_utf8_lossy(&buf[..n]).starts_with("GET /empty");
|
||||||
|
let body = if empty { "\n" } else {
|
||||||
|
"<?xml version=\"1.0\"?><rss version=\"2.0\"><channel><title>T</title><item><title>One</title><guid>g1</guid></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 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
|
||||||
|
|||||||
@@ -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());
|
||||||
|
|||||||
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>
|
||||||
@@ -259,7 +259,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 +270,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 +331,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