diff --git a/README.md b/README.md index 8f412b5..4a3c8c8 100644 --- a/README.md +++ b/README.md @@ -10,11 +10,12 @@ share for translations. One package, `session_ops`, deployed to one self-hosted | --- | --- | --- | | `zendesk-digest` | Weekday Discord digest of the Zendesk tickets awaiting a reply, after closing positive app-store reviews | [zendesk-digest](docs/jobs/zendesk-digest.md) | | `zendesk-relay` | Drafts and sends Zendesk replies from `claude:` private notes | [zendesk-relay](docs/jobs/zendesk-relay.md) | +| `session-ops-discord` | `/run` a job and `/mau-upload` an export from Discord | [session-ops-discord](docs/jobs/session-ops-discord.md) | | `github-prs-digest` | Weekday Discord digest of open pull requests from outside contributors | [github-prs-digest](docs/jobs/github-prs-digest.md) | | `crowdin-duplicates` | Crowdin string slots holding more than one translation, as they open and close | [crowdin-duplicates](docs/jobs/crowdin-duplicates.md) | | `crowdin-sync` | Weekday: Crowdin translations into iOS, Android and the localization module, and its submodule bumped in each client | [crowdin-sync](docs/jobs/crowdin-sync.md) | | `snode-list` | Weekday: the fallback service node list from the seed nodes into dynamic-assets, Desktop and iOS | [snode-list](docs/jobs/snode-list.md) | -| `release-stats` | On demand: download counts of the latest releases | [release-stats](docs/jobs/release-stats.md) | +| `mau` | Monthly: active users per platform and in total, from the store exports dropped in its inbox | [mau](docs/jobs/mau.md) | | `session-ops-silence` | Discord alerts for a job that failed, or stopped running | [session-ops-silence](docs/jobs/session-ops-silence.md) | | `token-expiry` | Discord alerts 14 days, 7 days and 24 hours before a token expires | [token-expiry](docs/jobs/token-expiry.md) | diff --git a/deploy/README.md b/deploy/README.md index a1c6cf3..6731ea0 100644 --- a/deploy/README.md +++ b/deploy/README.md @@ -7,8 +7,10 @@ does the work: accounts, venv, env files, units, timers, and migrating an older | Unit | What it is | | --- | --- | | `session-ops@.timer` → `.service` | One per job; a generated drop-in sets its account, env files and schedule. | +| `session-ops@.path` → `.service` | For a job with a `watch`: starts it when a matching file lands in its state directory. | | `session-ops-queue.timer` → `.service` | Starts the jobs in `jobs.toml`'s `[queue]`, which then run one at a time in its order. | | `zendesk-relay.service` | Always on, `127.0.0.1:8080`: Zendesk's `claude:` note webhooks. | +| `session-ops-discord.service` | Always on, `127.0.0.1:8081`: the `/run` and `/mau-upload` slash commands. See [Discord commands](#discord-commands). | | `session-ops-alert@.service` | Every unit's `OnFailure=` backstop; see [session-ops-silence](../docs/jobs/session-ops-silence.md). | `session-ops list` shows the jobs; `session-ops run [--dry-run] [-- job arguments]` @@ -52,6 +54,37 @@ From a host that ran the digests out of `/opt/zendesk`, `install.sh` copies thei files and state over. Remove `/opt/zendesk`, `/etc/zendesk` and `/var/lib/zendesk` (and the `github-prs` equivalents) once both digests have run from the new units. +## Discord commands + +What `/run` and `/mau-upload` do, and who may run them: +[session-ops-discord](../docs/jobs/session-ops-discord.md). To set them up: + +1. In the Developer Portal, create an application. Put its public key, its application + id, the server's id and who may run the commands in `/etc/session-ops/discord.env`, + then run `install.sh` again. +2. Invite it with the `applications.commands` scope only: + `https://discord.com/oauth2/authorize?client_id=&scope=applications.commands`. +3. Copy the `location = /discord/interactions` block from + [`nginx-webhooks.conf`](nginx-webhooks.conf) into the live + `/etc/nginx/sites-available/webhooks.session.codes`, which certbot owns, then + `nginx -t && systemctl reload nginx`. Until then the route is a 404 and Discord will + not save the URL. +4. Set its Interactions Endpoint URL to `https://webhooks.session.codes/discord/interactions`. + Discord checks it with a signed PING before saving. +5. Register the commands, with the bot token from the Bot page, typed in rather than kept: + + ```bash + read -rs DISCORD_BOT_TOKEN && export DISCORD_BOT_TOKEN + set -a && . /etc/session-ops/discord.env && set +a + /opt/session-ops/.venv/bin/session-ops discord-register + ``` + +6. The commands start hidden from everyone but server administrators. In Server + Settings → Integrations, allow them for the role or channel that runs jobs. Seeing + them is not running them: the relay's allowlist still decides. + +Register again after changing which jobs have `discord = true`. + ## Secrets Each `/etc/session-ops/.env` has a commented `.env.example` beside it, @@ -69,6 +102,8 @@ systemctl start session-ops@.service && journalctl -fu session-ops@ systemctl start session-ops-alert@test.service # posts to the alerts channel curl -sS -o /dev/null -w '%{http_code}\n' -X POST 127.0.0.1:8080/zendesk/notes \ -H 'Content-Type: application/json' -d '{"ticket_id":"1"}' # expect 401 +curl -sS -o /dev/null -w '%{http_code}\n' -X POST 127.0.0.1:8081/discord/interactions -d '{}' # expect 401 +runuser -u opsbot -- systemctl --no-ask-password start session-ops@token-expiry.service # expect Access denied ``` ## Rehearsing on a spare host @@ -85,7 +120,10 @@ To end it, stop the timers, then close the rehearsal pull requests: for link in /etc/systemd/system/timers.target.wants/session-ops@*.timer; do [ -L "$link" ] && systemctl disable --now "${link##*/}" done -systemctl disable --now session-ops-queue.timer zendesk-relay.service +for link in /etc/systemd/system/paths.target.wants/session-ops@*.path; do + [ -L "$link" ] && systemctl disable --now "${link##*/}" +done +systemctl disable --now session-ops-queue.timer zendesk-relay.service session-ops-discord.service for repo in session-android session-ios session-localization session-desktop-dynamic-assets \ session-desktop session-app session-website session-appium session-playwright; do gh pr list -R "session-foundation/$repo" --state open --json number,headRefName \ diff --git a/deploy/env/alerts.env.example b/deploy/env/alerts.env.example index 3a0d509..d36b964 100644 --- a/deploy/env/alerts.env.example +++ b/deploy/env/alerts.env.example @@ -1,3 +1,3 @@ -# /etc/session-ops/alerts.env: the OnFailure backstop, the silence checker and -# release-stats. Empty disables the first two, and install.sh warns until it is set. +# /etc/session-ops/alerts.env: the OnFailure backstop and the silence checker. +# Empty disables both, and install.sh warns until it is set. ALERT_DISCORD_WEBHOOK_URL= diff --git a/deploy/env/discord.env.example b/deploy/env/discord.env.example new file mode 100644 index 0000000..4d064eb --- /dev/null +++ b/deploy/env/discord.env.example @@ -0,0 +1,14 @@ +# /etc/session-ops/discord.env: the Discord app behind /run and /mau-upload. +# A `#` starts a comment only as a line's first character. Restart session-ops-discord +# after an edit: it reads this file at start. + +# General Information in the Developer Portal. Unset refuses every interaction. +DISCORD_PUBLIC_KEY= +# The one server the commands are registered on and answered from. +DISCORD_GUILD_ID= +# Comma-separated; either list grants. Both empty refuses everybody. +ALLOWED_USER_IDS= +ALLOWED_ROLE_IDS= +# For `session-ops discord-register`. Its DISCORD_BOT_TOKEN is typed in when registering, +# not kept here: the relay never needs it. +DISCORD_APP_ID= diff --git a/deploy/env/mau.env.example b/deploy/env/mau.env.example new file mode 100644 index 0000000..4026ab3 --- /dev/null +++ b/deploy/env/mau.env.example @@ -0,0 +1,7 @@ +# /etc/session-ops/mau.env: the monthly active users post. +# Where the figures, the reminders and the job's own failures are posted. +MAU_DISCORD_WEBHOOK_URL= +# An App Store Connect team API key with the Sales and Reports role, for Apple's opt-in +# rate: its issuer and key IDs here, its .p8 at /etc/session-ops/asc-key.p8 (mode 600). +ASC_ISSUER_ID= +ASC_KEY_ID= diff --git a/deploy/install.sh b/deploy/install.sh index 74d4fde..f8c3a8f 100755 --- a/deploy/install.sh +++ b/deploy/install.sh @@ -37,7 +37,7 @@ account() { id -u "$1" >/dev/null 2>&1 || useradd --system --no-create-home --home /nonexistent --shell /usr/sbin/nologin "$1" } -for user in ghdigest crowdin publisher sessionops; do account "$user"; done +for user in ghdigest crowdin publisher sessionops opsbot; do account "$user"; done # The Claude Code CLI keeps its binary, login and cache under this account's $HOME. if ! id -u zendesk >/dev/null 2>&1; then useradd --system --home /home/zendesk --shell /usr/sbin/nologin zendesk @@ -68,15 +68,16 @@ move /etc/github-prs/env "$ETC/github-prs.env" 600 root move "$ETC/env" "$ETC/alerts.env" 600 root # systemd reads EnvironmentFile= as root before dropping privileges, so these need # no group: a job's account cannot read another job's secrets. -for name in zendesk github-prs crowdin publish alerts; do +for name in zendesk github-prs crowdin publish alerts mau discord; do [ -e "$ETC/$name.env" ] || install -m 600 /dev/null "$ETC/$name.env" done # What each file takes, commented; the env files stay empty until filled, since a job # is enabled once its env files have content. install -m 644 "$ROOT"/deploy/env/*.example "$ETC/" -# The publishing units load this as a credential, and a missing file would stop them -# starting; left empty, their runs exit naming the key. +# The publishing units and mau load these as credentials, and a missing file would stop +# them starting; left empty, their runs fail naming the key. [ -e "$ETC/github-app.pem" ] || install -m 600 /dev/null "$ETC/github-app.pem" +[ -e "$ETC/asc-key.p8" ] || install -m 600 /dev/null "$ETC/asc-key.p8" install -d -m 755 "$STATE" install -d -o ghdigest -g ghdigest "$STATE/github-prs-digest" @@ -92,6 +93,17 @@ if [ -e "$NEW_HOUSE" ] && grep -qx "ZENDESK_HOUSE_ANSWERS=$OLD_HOUSE" "$ETC/zend echo "pointed ZENDESK_HOUSE_ANSWERS in $ETC/zendesk.env at $NEW_HOUSE" fi +# A watched job's inbox: root drops files in, and the job's account moves them out. +# mau's also takes /mau-upload's, from the Discord relay's account. +"$OPS" list --watched | while read -r job user dir; do + install -d -o "$user" -g "$user" -m 711 "$STATE/$job" + if [ "$job" = mau ]; then + install -d -o "$user" -g opsbot -m 770 "$dir" + else + install -d -o "$user" -g "$user" -m 700 "$dir" + fi +done + install -m 644 "$ROOT/deploy/session-ops.tmpfiles" /etc/tmpfiles.d/session-ops.conf systemd-tmpfiles --create session-ops.conf @@ -104,14 +116,21 @@ rm -f "$UNITS/zendesk-alert@.service" "$UNITS/github-prs-alert@.service" systemctl disable --now crowdin-relay.service 2>/dev/null || true rm -f "$UNITS/crowdin-relay.service" -install -m 644 "$ROOT"/deploy/*.service "$ROOT"/deploy/*.timer "$UNITS/" +install -m 644 "$ROOT"/deploy/*.service "$ROOT"/deploy/*.timer "$ROOT"/deploy/*.path "$UNITS/" # Only the generated files go, so a drop-in added by hand survives. rm -f "$UNITS"/session-ops@*.service.d/job.conf "$UNITS"/session-ops@*.timer.d/schedule.conf \ - "$UNITS"/session-ops-queue.timer.d/schedule.conf + "$UNITS"/session-ops@*.path.d/watch.conf "$UNITS"/session-ops-queue.timer.d/schedule.conf "$OPS" units --out "$UNITS" >/dev/null +# polkit reloads its rules when this changes. +if [ -d /etc/polkit-1/rules.d ]; then + "$OPS" polkit --out /etc/polkit-1/rules.d/50-session-ops-discord.rules >/dev/null +else + echo "polkit is missing, so /run in Discord can start no job" >&2 +fi READY=$("$OPS" list --ready) QUEUED=$("$OPS" list --queued) +WATCHED=$("$OPS" list --watched | cut -d' ' -f1) listed() { printf '%s\n' $2 | grep -qxF "$1"; } # What the queue's timer starts: its ready jobs, rebuilt from scratch each install. WANTS="$UNITS/session-ops-queue.service.wants" @@ -142,6 +161,21 @@ for job in $READY; do echo "enabled session-ops@$job.timer" fi done +for link in "$UNITS"/paths.target.wants/session-ops@*.path; do + [ -L "$link" ] || continue + job=${link##*/session-ops@} + job=${job%.path} + if ! listed "$job" "$READY" || ! listed "$job" "$WATCHED"; then + systemctl disable --now "session-ops@$job.path" >/dev/null + echo "disabled session-ops@$job.path (no longer a ready job with a watch)" + fi +done +for job in $WATCHED; do + if listed "$job" "$READY"; then + systemctl enable --now "session-ops@$job.path" >/dev/null + echo "enabled session-ops@$job.path" + fi +done if [ -d "$WANTS" ]; then systemctl enable --now session-ops-queue.timer >/dev/null echo "enabled session-ops-queue.timer" @@ -151,16 +185,21 @@ fi for job in $("$OPS" list --not-ready); do echo "not enabled: session-ops@$job (its env file is empty)" done -if [ -s "$ETC/zendesk.env" ]; then - systemctl enable zendesk-relay.service >/dev/null - # A relay that hit its start limit refuses `start` until the limit is cleared. - systemctl reset-failed zendesk-relay.service 2>/dev/null || true - systemctl try-restart zendesk-relay.service - systemctl start zendesk-relay.service - echo "running zendesk-relay.service" -else - echo "not enabled: zendesk-relay.service ($ETC/zendesk.env is empty)" -fi +# An always-on service, enabled once its env file has content. +relay() { + if [ -s "$ETC/$2.env" ]; then + systemctl enable "$1.service" >/dev/null + # A relay that hit its start limit refuses `start` until the limit is cleared. + systemctl reset-failed "$1.service" 2>/dev/null || true + systemctl try-restart "$1.service" + systemctl start "$1.service" + echo "running $1.service" + else + echo "not enabled: $1.service ($ETC/$2.env is empty)" + fi +} +relay zendesk-relay zendesk +relay session-ops-discord discord if [ ! -s "$ETC/alerts.env" ]; then cat >&2 <:/var/lib/session-ops/mau/inbox/ +``` + +`rsync`, not `scp`: it writes to a hidden temporary name and renames it once complete, +and the job only reads `*.csv`. App Store Connect's Active Devices export is refused: +it counts each day apart, and the App Store Connect API's App Sessions report is no +substitute, since summing its rows counts a device once per app version, OS and +territory it used in the month, 5% to 20% too many. + +## What a run does + +1. Merges every export in `inbox/` into `history.json`, one figure per platform and day, + then moves it to `done/`. An export may cover any range; where two give a day different + figures, the later one wins, and the change is kept until the month posts. +2. Moves a file it cannot read to `rejected/` and fails the run naming it, after the + rest of the run. +3. Posts last month once `history.json` holds a final figure for its last day on both + platforms. Play's figures trail by about eight days, so that is around the 9th. Apple + gives a day it is still counting as a fraction, which does not count as final. Until + then, from the 10th, each run posts a reminder naming the platforms still missing. +4. Lists in the post the changes to the two month ends it compares, records the month as + posted, so neither the timer nor a later export posts it again, and forgets the + changes. + +Should a run fail before the inbox is filed, say on an unreadable `history.json`, it +moves every export in `inbox/` to `rejected/` first: the path unit starts the job again +for as long as a file matches, and gives up watching after five starts. Once the cause +is fixed, move them back into `inbox/`. + +`history.json` cannot be rebuilt by a re-run: an unreadable one stops the job. If it is +lost, drop an export covering the last 365 days for each platform. + +## Desktop downloads + +The latest Desktop release that is neither a draft nor a pre-release. Updates count +where the updater fetches an installer. + +| Platform | Counted | +| --- | --- | +| Linux | `.deb`, `.AppImage`, `.rpm`, `.freebsd`, and Flathub's installs of `network.loki.Session` since the release day | +| macOS | `.dmg`, `.zip` | +| Windows | `.exe` | + +`.blockmap`, `latest*.yml` and `signature.asc` are left out. Flathub builds from the +GitHub `.deb` once, so its installs are not in GitHub's counts; Homebrew's `session` +cask downloads the GitHub `.dmg`, so it already is. Flathub keeps 180 days of daily +installs, so the run fails on a release older than that rather than undercount. + +## Apple's opt-in rate + +The job reads the App Opt In report from every analytics report request on Session's +App Store Connect app: the one-time snapshot for history, and the ongoing one for each +new day. A Sales and Reports key may read them but not create them; if Apple stops the +ongoing request for inactivity, the run fails saying so, and an Admin key has to +create a new one (`POST /v1/analyticsReportRequests`, `accessType: ONGOING`). + +## Android outside Play + +The `.apk` files of the latest Android release on GitHub that is neither a draft nor a +pre-release, since its release, updates included. F-Droid builds its own APK and +publishes no download counts, so its figure is 10% of every other figure in the post. diff --git a/docs/jobs/release-stats.md b/docs/jobs/release-stats.md deleted file mode 100644 index 3c41d60..0000000 --- a/docs/jobs/release-stats.md +++ /dev/null @@ -1,27 +0,0 @@ -# Release Download Statistics - -Download counts of the last ten Desktop and Android releases, per installer, from the -public releases API, as two CSV files. - -| | | -| --- | --- | -| Runs | on demand only: `systemctl start session-ops@release-stats.service` | -| Secrets | none; `/etc/session-ops/alerts.env` for its failures | -| Dry run | `session-ops run release-stats --dry-run` prints the CSVs and writes nothing | -| Logs | `journalctl -u session-ops@release-stats -n 50 --no-pager`, which carries both CSVs | - -The files land in `/var/lib/session-ops/release-stats/runs//` and are pruned -after 14 days. From a checkout, `uv run python -m session_ops.platforms.release_stats ---out .` writes them to a timestamped folder in the current directory, `.//`. - -| Desktop column | Assets counted | -| --- | --- | -| `.deb`, `.appimage`, `.rpm`, `.exe` | by extension | -| `.dmg_arm64`, `.dmg_x64`, `.zip_arm64`, `.zip_x64` | by extension and architecture | - -| Android column | Assets counted | -| --- | --- | -| `.aab` | `play-release` bundles | -| `.apk_arm64`, `.apk_armv7a`, `.apk_x86_64` | `play-release` APKs for `arm64-v8a`, `armeabi-v7a`, `x86_64` | -| `.apk_x86` | `play-release` x86 APKs, not x86_64 | -| `.apk_universal_play`, `.apk_universal_huawei` | universal APKs per store | diff --git a/docs/jobs/session-ops-discord.md b/docs/jobs/session-ops-discord.md new file mode 100644 index 0000000..55dc9af --- /dev/null +++ b/docs/jobs/session-ops-discord.md @@ -0,0 +1,29 @@ +# Discord Commands for the Jobs + +| | | +| --- | --- | +| Runs | `session-ops-discord.service`, always on, behind nginx at `POST /discord/interactions` | +| Secrets | `/etc/session-ops/discord.env`: the app's public key, the server and who may run the commands | +| Setup | [deploy/README.md](../../deploy/README.md#discord-commands) | +| Logs | `journalctl -u session-ops-discord -n 50 --no-pager`: who ran what, and each outcome | + +| Command | Does | +| --- | --- | +| `/run job:` | Starts `session-ops@.service` for a job `jobs.toml` marks `discord = true`. Refused while that job or any queued job is running or waiting to, or when the queue's next run is sooner than the job's timeout from now. The job posts its outcome in its own channel. | +| `/mau-upload file:` | Downloads the attachment, checks it parses as a Play Console or App Store Connect export, and moves it into [mau](mau.md)'s inbox, whose path unit runs the job. A file it would reject is refused in Discord and never reaches the inbox. | + +The app has no gateway connection and only the `applications.commands` scope, so Discord +sends it the commands run against it and nothing else: no messages, members or other +channels. Each request carries who ran it, their roles, the server, the channel and the +options chosen. + +A request is refused unless Discord signed it in the last five minutes, it comes from +`DISCORD_GUILD_ID`, and its author is in `ALLOWED_USER_IDS` or holds a role in +`ALLOWED_ROLE_IDS`. Both lists empty refuses everybody. The commands are registered +hidden from everyone but server administrators, and the server's Integrations settings +show them to the right role; that only hides them, and these checks are what refuse +everyone else, administrators included. + +The relay runs as `opsbot`. The polkit rule `install.sh` writes from `jobs.toml` lets that +account start the `discord = true` jobs and nothing else, and the mau inbox is the one +path it can write. diff --git a/src/session_ops/github_prs/digest.py b/src/session_ops/github_prs/digest.py index f2247b6..a75a7d5 100755 --- a/src/session_ops/github_prs/digest.py +++ b/src/session_ops/github_prs/digest.py @@ -236,7 +236,7 @@ def window_start(state, now, window_hours, retention_days=DEFAULT_RETENTION_DAYS the state says reporting is complete. A fixed window alone drops whatever moved in a gap longer than it: a DST weekend, - a timer's random delay, a host down across a run, a run whose post failed partway. + a host down across a run, a run whose post failed partway. Never further back than the retention, past which the state has forgotten what it reported anyway. """ diff --git a/src/session_ops/jobs.toml b/src/session_ops/jobs.toml index fdaddc7..88dc380 100644 --- a/src/session_ops/jobs.toml +++ b/src/session_ops/jobs.toml @@ -10,6 +10,9 @@ # channel_env where it posts, and where its own failures are reported; a job # without one reports to ALERT_DISCORD_WEBHOOK_URL # unit further [Service] lines, over the template's hardening +# discord offered to `/run` in Discord, and startable by the relay's account +# watch a glob under {state}: a file matching it starts the job, which must +# move it out, or the path unit starts it again # # max_age_hours for a queued job, its queue's: the weekend's 72 h (73 h across # April's DST change), plus every job ahead of it running to its timeout @@ -71,6 +74,7 @@ channel_env = "CROWDIN_DISCORD_WEBHOOK_URL" max_age_hours = 80 # About 9 minutes; a full scan without --croql is ~110k requests, about 85 minutes. timeout = "2h" +discord = true [[job]] name = "session-ops-silence" @@ -111,6 +115,7 @@ env = ["CROWDIN_API_TOKEN", "PUBLISH_GIT_AUTHOR"] channel_env = "CROWDIN_DISCORD_WEBHOOK_URL" max_age_hours = 80 timeout = "1h" +discord = true # Readable by this unit alone, at $CREDENTIALS_DIRECTORY; publishing needs it. unit = ["LoadCredential=github-app.pem:/etc/session-ops/github-app.pem"] @@ -127,10 +132,19 @@ timeout = "10min" unit = ["LoadCredential=github-app.pem:/etc/session-ops/github-app.pem"] [[job]] -name = "release-stats" -description = "Download counts of the latest Desktop and Android releases" -entry = "session_ops.platforms.release_stats:main" -args = ["--out", "{state}/runs"] +name = "mau" +description = "Monthly active users per platform, posted once a month" +entry = "session_ops.platforms.mau:main" +args = ["--state", "{state}"] user = "sessionops" -env_files = ["/etc/session-ops/alerts.env"] +env_files = ["/etc/session-ops/mau.env"] +env = ["MAU_DISCORD_WEBHOOK_URL", "ASC_ISSUER_ID", "ASC_KEY_ID"] +channel_env = "MAU_DISCORD_WEBHOOK_URL" +# An hour clear of the queue's posts at 10:00. +schedule = "*-*-10 11:00 Australia/Melbourne" +watch = "inbox/*.csv" +# The longest month, plus the hour a DST change adds and slack for the run. +max_age_hours = 750 timeout = "5min" +# The App Store Connect API key, for Apple's opt-in rate. +unit = ["LoadCredential=asc-key.p8:/etc/session-ops/asc-key.p8"] diff --git a/src/session_ops/ops/discord_commands.py b/src/session_ops/ops/discord_commands.py new file mode 100644 index 0000000..bb4c690 --- /dev/null +++ b/src/session_ops/ops/discord_commands.py @@ -0,0 +1,63 @@ +"""The slash commands session-ops-discord.service answers, and their registration. + + session-ops discord-register # with discord.env's variables in the environment + +Registered on one guild rather than globally: a guild's commands update at once, and +nobody outside that server can see or run them. They start hidden from every member +but server administrators; the server's Integrations settings grant them to a role or +channel. That only hides them, and the relay's allowlist is what refuses everyone else, +administrators included. +""" +import os +import sys + +from session_ops.ops import registry +from session_ops.shared import discord, http + +API = "https://discord.com/api/v10" +RUN = "run" +MAU_UPLOAD = "mau-upload" +MAU_JOB = "mau" + +OPTION_STRING = 3 +OPTION_ATTACHMENT = 11 +CHAT_INPUT = 1 +# Discord's cap on a choice's label. +MAX_CHOICE_NAME = 100 + + +def commands(jobs): + choices = [{"name": discord.clip(f"{job.name}: {job.description}", MAX_CHOICE_NAME), + "value": job.name} for job in jobs if job.discord] + hidden = {"type": CHAT_INPUT, "default_member_permissions": "0"} + return [ + {**hidden, "name": RUN, "description": "Start a session-ops job now", + "options": [{"type": OPTION_STRING, "name": "job", "required": True, + "description": "The job to start", "choices": choices}]}, + {**hidden, "name": MAU_UPLOAD, + "description": "Hand the mau job a Play Console or App Store Connect export", + "options": [{"type": OPTION_ATTACHMENT, "name": "file", "required": True, + "description": "Play's MAU report, or App Store Connect's Active in " + "Last 30 Days, as CSV"}]}, + ] + + +def register(session, app_id, guild_id, token, jobs): + """Replace every command on the guild with these. Returns the names registered.""" + response = session.request("PUT", f"{API}/applications/{app_id}/guilds/{guild_id}/commands", + headers={"Authorization": f"Bot {token}"}, json=commands(jobs)) + if response.status_code != 200: + raise SystemExit(f"Discord refused the commands: {response.status_code} " + f"{response.text[:300]}") + return [command["name"] for command in response.json()] + + +def main(): + missing = [name for name in ("DISCORD_APP_ID", "DISCORD_GUILD_ID", "DISCORD_BOT_TOKEN") + if not os.environ.get(name)] + if missing: + sys.exit(f"missing {', '.join(missing)} in the environment") + names = register(http.Session(attempts=3), os.environ["DISCORD_APP_ID"], + os.environ["DISCORD_GUILD_ID"], os.environ["DISCORD_BOT_TOKEN"], + registry.load()) + print("Registered /" + ", /".join(names)) diff --git a/src/session_ops/ops/discord_relay.py b/src/session_ops/ops/discord_relay.py new file mode 100644 index 0000000..6993772 --- /dev/null +++ b/src/session_ops/ops/discord_relay.py @@ -0,0 +1,317 @@ +"""Discord slash commands for session-ops jobs: POST /discord/interactions. + + /run job: starts session-ops@.service + /mau-upload file: puts a Play Console or App Store Connect export in mau's + inbox, whose path unit then runs the job + +Discord signs every interaction with the app's Ed25519 key. One is refused unless it +comes from DISCORD_GUILD_ID and from someone in ALLOWED_USER_IDS or ALLOWED_ROLE_IDS. +The account this runs as may start only the jobs jobs.toml marks `discord`: the polkit +rule `session-ops polkit` writes refuses it everything else. + +Config (env vars, /etc/session-ops/discord.env): + DISCORD_PUBLIC_KEY the app's public key. Unset refuses every request + DISCORD_GUILD_ID the one server it answers + ALLOWED_USER_IDS comma-separated Discord user ids + ALLOWED_ROLE_IDS comma-separated role ids + +Usage: + uvicorn session_ops.ops.discord_relay:app --host 127.0.0.1 --port 8081 +""" +import json +import os +import subprocess +import threading +import time +from urllib.parse import urlsplit + +from fastapi import BackgroundTasks, FastAPI, Request, Response +from nacl.exceptions import BadSignatureError +from nacl.signing import VerifyKey +from starlette.concurrency import run_in_threadpool + +from session_ops.ops import registry +from session_ops.ops.discord_commands import API, MAU_JOB, MAU_UPLOAD, RUN +from session_ops.platforms import mau +from session_ops.shared import http + +# Bounds replay of a request Discord genuinely signed; the signature covers the timestamp. +MAX_SIGNATURE_AGE_SECONDS = 300 +# Without it, a host whose clock trails Discord's refuses every interaction, the +# endpoint-registering PING included. +MAX_CLOCK_SKEW_SECONDS = 60 +# A year of daily figures is a few kilobytes. +MAX_UPLOAD_BYTES = 1 << 20 +ATTACHMENT_HOSTS = frozenset({"cdn.discordapp.com", "media.discordapp.net"}) +SYSTEMCTL_TIMEOUT_SECONDS = 10 +RUNNING_STATES = frozenset({"active", "activating", "deactivating", "reloading"}) +USEC_INFINITY = 2**64 - 1 + +INTERACTION_PING = 1 +INTERACTION_COMMAND = 2 +RESPONSE_PONG = 1 +RESPONSE_MESSAGE = 4 +RESPONSE_DEFERRED = 5 +EPHEMERAL = 64 +NO_PINGS = {"parse": []} + +app = FastAPI(docs_url=None, redoc_url=None, openapi_url=None) +# Across the check and the start, so two /run a moment apart cannot both pass; one process only. +_starting = threading.Lock() + + +def env(name, default=None): + return os.environ.get(name, default) + + +def id_list(raw): + return [item.strip() for item in (raw or "").split(",") if item.strip()] + + +def signature_ok(raw, signature, timestamp, key): + if not (signature and timestamp and key): + return False + try: + age = time.time() - float(timestamp) + except (TypeError, ValueError): + return False + if age < -MAX_CLOCK_SKEW_SECONDS or age > MAX_SIGNATURE_AGE_SECONDS: + return False + try: + VerifyKey(bytes.fromhex(key)).verify(timestamp.encode() + raw, bytes.fromhex(signature)) + except (BadSignatureError, ValueError): + return False + return True + + +def user_id(interaction): + member = interaction.get("member") or {} + return (member.get("user") or interaction.get("user") or {}).get("id") + + +def refusal(interaction): + """Why this person may not use the commands here, or None. + + Both allowlists empty refuses everybody: an unconfigured relay must not mean an open one. + A DM carries no guild_id, so it is refused with the wrong server. + """ + guild = env("DISCORD_GUILD_ID") + if not guild or interaction.get("guild_id") != guild: + return "These commands do not work here." + users = id_list(env("ALLOWED_USER_IDS")) + roles = id_list(env("ALLOWED_ROLE_IDS")) + if not users and not roles: + return "Neither ALLOWED_USER_IDS nor ALLOWED_ROLE_IDS is set on the relay." + if user_id(interaction) in users: + return None + if any(role in roles for role in (interaction.get("member") or {}).get("roles") or []): + return None + return "You are not on the list of people who can run session-ops jobs." + + +def reply(text, ephemeral=True): + data = {"content": text, "allowed_mentions": NO_PINGS} + if ephemeral: + data["flags"] = EPHEMERAL + return {"type": RESPONSE_MESSAGE, "data": data} + + +def option(interaction, name): + for item in (interaction.get("data") or {}).get("options") or []: + if item.get("name") == name: + return item.get("value") + return None + + +def systemctl(*args): + return subprocess.run(["systemctl", "--no-ask-password", *args], check=False, + capture_output=True, text=True, timeout=SYSTEMCTL_TIMEOUT_SECONDS) + + +def unit_name(job_name): + return f"session-ops@{job_name}.service" + + +def busy(job): + """The first of `job` and the queue's jobs that is running or waiting to, else None. + + The queue runs its jobs one at a time so no two share the host or their APIs; a run + started from here keeps to that. + """ + units = [unit_name(name) for name in dict.fromkeys([job.name, *registry.load_queue().jobs])] + waiting = {line.split()[1] for line in + systemctl("list-jobs", "--no-legend", "--plain").stdout.splitlines() + if len(line.split()) > 1} + states = systemctl("is-active", *units).stdout.split() + if len(states) != len(units): + raise RuntimeError(f"systemctl is-active gave {len(states)} states for {len(units)} units") + for unit, state in zip(units, states): + if unit in waiting or state in RUNNING_STATES: + return unit + return None + + +def next_queue_run(): + """The queue timer's next run as a Unix time, or None while it has none. + + From list-timers' JSON, in microseconds: systemd 252 prints this property as local + time under `show`, --timestamp=unix or not. + """ + timers = json.loads(systemctl("list-timers", "--all", "--output=json", + registry.QUEUE_TIMER).stdout or "[]") + next_us = timers[0].get("next") if timers else None + return next_us / 1e6 if next_us and next_us != USEC_INFINITY else None + + +def hold_off(job): + """Why `job` cannot start now, or None.""" + blocking = busy(job) + if blocking: + return f"⏳ `{blocking}` is running or waiting to; try again once it has finished." + queue_at = next_queue_run() + # A run is killed at its timeout, so the queue starting later cannot overlap it. + if queue_at is None or queue_at >= time.time() + job.timeout_seconds: + return None + at = f" ()" + if job.queued: + return (f"⏳ The queue runs **{job.name}** at {at}, sooner than a run started now " + f"would be sure to end.") + return (f"⏳ The queue starts at {at}, sooner than **{job.name}** would be sure to end; " + f"try again once it has run.") + + +def handle_run(interaction): + name = option(interaction, "job") + job = next((job for job in registry.load() if job.discord and job.name == name), None) + if job is None: + return reply(f"❌ `{name}` is not a job /run can start.") + unit = unit_name(job.name) + try: + with _starting: + held = hold_off(job) + started = None if held else systemctl("start", "--no-block", unit) + except (OSError, subprocess.SubprocessError, RuntimeError, ValueError) as exc: + print(f"/run {job.name}: {exc!r}", flush=True) + return reply("❌ Could not ask systemd on the host; see " + "`journalctl -u session-ops-discord`.") + if held: + return reply(held) + if started.returncode != 0: + error = started.stderr.strip() + print(f"/run {job.name} by {user_id(interaction)}: {error}", flush=True) + return reply(f"❌ systemd did not start **{job.name}**: {error[:300]}") + print(f"/run {job.name} by {user_id(interaction)}: started", flush=True) + return reply(f"▶️ <@{user_id(interaction)}> started **{job.name}**. It posts in its own " + f"channel when it is done.", ephemeral=False) + + +def attachment_problem(attachment): + if not attachment: + return "Discord sent no file with the command." + filename = attachment.get("filename") or "" + if not filename.lower().endswith(".csv"): + return f"Both stores export a .csv, and this is `{filename}`." + if not isinstance(attachment.get("size"), int) or attachment["size"] > MAX_UPLOAD_BYTES: + return f"`{filename}` is larger than any MAU export, at {attachment.get('size')} bytes." + url = urlsplit(attachment.get("url") or "") + if url.scheme != "https" or url.hostname not in ATTACHMENT_HOSTS: + return "The file is not on Discord's CDN." + return None + + +def inbox(): + return os.path.dirname(registry.get(MAU_JOB).watch_glob) + + +def place_export(session, attachment, name, uploader): + """Download the export into mau's inbox if it parses. Returns the message for Discord.""" + response = session.request("GET", attachment["url"]) + if response.status_code != 200: + return f"❌ Discord's CDN answered {response.status_code} for the file; send it again." + if len(response.content) > MAX_UPLOAD_BYTES: + return "❌ The file Discord served is larger than any MAU export." + folder = inbox() + # mau's glob skips dotfiles, so its path unit cannot start on a file still being checked. + hidden = os.path.join(folder, f".{name}") + try: + with open(hidden, "wb") as handle: + handle.write(response.content) + # mau runs as another account; the inbox's 0770 keeps everyone else out. + os.fchmod(handle.fileno(), 0o644) + try: + platform, days = mau.parse_export(hidden) + except mau.Rejected as exc: + return f"❌ That is not the export mau reads: {exc}" + os.replace(hidden, os.path.join(folder, name)) + finally: + if os.path.exists(hidden): + os.unlink(hidden) + return (f"📥 <@{uploader}>'s {mau.LABELS[platform]} export, {len(days)} " + f"day{'s' if len(days) != 1 else ''} from " + f"{min(days)} to {max(days)}, is in mau's inbox, and the job is running on it.") + + +def edit_original(session, interaction, text): + url = (f"{API}/webhooks/{interaction['application_id']}/{interaction['token']}" + f"/messages/@original") + try: + session.request("PATCH", url, attempts=2, + json={"content": text, "allowed_mentions": NO_PINGS}) + except Exception as exc: # noqa: BLE001 — the outcome is in the journal either way + print(f"could not tell Discord: {exc!r}", flush=True) + + +def drop_export(interaction, attachment): + """Never raises: Discord shows "thinking…" until this edits that message.""" + session = http.Session(attempts=3, timeout=30) + uploader = user_id(interaction) + try: + outcome = place_export(session, attachment, f"discord-{interaction['id']}.csv", uploader) + except Exception as exc: # noqa: BLE001 — reported to Discord below + print(f"/mau-upload by {uploader}: {exc!r}", flush=True) + outcome = "❌ The upload failed on the host; see `journalctl -u session-ops-discord`." + print(f"/mau-upload by {uploader}: {outcome}", flush=True) + edit_original(session, interaction, outcome) + + +def handle_mau_upload(interaction, background): + resolved = ((interaction.get("data") or {}).get("resolved") or {}).get("attachments") or {} + attachment = resolved.get(str(option(interaction, "file"))) + problem = attachment_problem(attachment) + if problem: + return reply(f"❌ {problem}") + background.add_task(run_in_threadpool, drop_export, interaction, attachment) + return {"type": RESPONSE_DEFERRED} + + +@app.get("/healthz") +def healthz(): + return {"ok": True} + + +@app.post("/discord/interactions") +async def interactions(request: Request, background: BackgroundTasks): + raw = await request.body() + if not signature_ok(raw, request.headers.get("x-signature-ed25519"), + request.headers.get("x-signature-timestamp"), env("DISCORD_PUBLIC_KEY")): + return Response("bad signature", status_code=401) + try: + interaction = json.loads(raw) + except ValueError: + interaction = None + if not isinstance(interaction, dict): + return Response("not an interaction", status_code=400) + kind = interaction.get("type") + if kind == INTERACTION_PING: + return {"type": RESPONSE_PONG} + if kind != INTERACTION_COMMAND: + return Response("unsupported interaction", status_code=400) + denied = refusal(interaction) + if denied: + return reply(f"❌ {denied}") + command = (interaction.get("data") or {}).get("name") + if command == RUN: + return await run_in_threadpool(handle_run, interaction) + if command == MAU_UPLOAD: + return handle_mau_upload(interaction, background) + return reply(f"❌ This relay has no /{command}.") diff --git a/src/session_ops/ops/registry.py b/src/session_ops/ops/registry.py index 9ed9186..daa388a 100644 --- a/src/session_ops/ops/registry.py +++ b/src/session_ops/ops/registry.py @@ -4,6 +4,7 @@ silence checker, so a job is added in one place. """ import os +import re import tomllib from dataclasses import dataclass, field, replace @@ -11,6 +12,30 @@ "jobs.toml") STATE_ROOT = "/var/lib/session-ops" QUEUE_TIMER = "session-ops-queue.timer" +# Discord's cap on a string option's choices, which `/run` offers these jobs as. +MAX_DISCORD_JOBS = 25 +SPAN_UNITS = { + **dict.fromkeys(("us", "usec"), 1e-6), **dict.fromkeys(("ms", "msec"), 1e-3), + **dict.fromkeys(("", "s", "sec", "second", "seconds"), 1), + **dict.fromkeys(("m", "min", "minute", "minutes"), 60), + **dict.fromkeys(("h", "hr", "hour", "hours"), 3600), + **dict.fromkeys(("d", "day", "days"), 86400), + **dict.fromkeys(("w", "week", "weeks"), 604800), +} +SPAN_PART = re.compile(r"\s*(\d+(?:\.\d+)?)\s*([a-z]*)") + + +def span_seconds(text): + """Seconds in a systemd time span such as "1h 30min"; a bare number is seconds.""" + parts = list(SPAN_PART.finditer(text)) + if not parts or "".join(part.group(0) for part in parts).strip() != text.strip(): + raise ValueError(f"{text!r} is not a time span") + total = 0 + for part in parts: + if part.group(2) not in SPAN_UNITS: + raise ValueError(f"{text!r} has an unknown unit {part.group(2)!r}") + total += float(part.group(1)) * SPAN_UNITS[part.group(2)] + return total @dataclass(frozen=True) @@ -30,6 +55,8 @@ class Job: unit: tuple = field(default=()) queued: bool = False after: tuple = () + watch: str = None + discord: bool = False @property def scheduled(self): @@ -39,10 +66,18 @@ def scheduled(self): def timer(self): return QUEUE_TIMER if self.queued else f"session-ops@{self.name}.timer" + @property + def timeout_seconds(self): + return span_seconds(self.timeout) + @property def state_dir(self): return os.path.join(STATE_ROOT, self.name) + @property + def watch_glob(self): + return os.path.join(self.state_dir, self.watch) if self.watch else None + def argv(self, dry_run=False, state_dir=None): """The job's arguments, with {state} standing for its state directory.""" state = state_dir or self.state_dir @@ -84,8 +119,18 @@ def load(path=REGISTRY): for k, v in row.items()})) jobs = _apply_queue(path, jobs, data.get("queue")) for job in jobs: + if job.watch and (os.path.isabs(job.watch) or ".." in job.watch.split("/")): + raise ValueError(f"{path}: {job.name}'s watch must stay inside its state directory") if job.scheduled and not isinstance(job.max_age_hours, (int, float)): raise ValueError(f"{path}: scheduled job {job.name} needs a numeric max_age_hours") + if not isinstance(job.discord, bool): + raise ValueError(f"{path}: {job.name}'s discord must be true or false") + try: + job.timeout_seconds + except (TypeError, ValueError) as exc: + raise ValueError(f"{path}: {job.name}'s timeout: {exc}") from None + if sum(job.discord for job in jobs) > MAX_DISCORD_JOBS: + raise ValueError(f"{path}: Discord offers at most {MAX_DISCORD_JOBS} jobs to /run") return jobs diff --git a/src/session_ops/ops/runner.py b/src/session_ops/ops/runner.py index f4ab382..731e6c9 100644 --- a/src/session_ops/ops/runner.py +++ b/src/session_ops/ops/runner.py @@ -3,6 +3,8 @@ session-ops list session-ops run github-prs-digest [--dry-run] [-- further job arguments] session-ops units --out /etc/systemd/system + session-ops polkit --out /etc/polkit-1/rules.d/50-session-ops-discord.rules + session-ops discord-register A run owns what every job would otherwise repeat: checking the environment it needs, a scratch directory, and telling Discord when it fails. The alert names the job, the @@ -235,12 +237,19 @@ def main(argv=None): help="Only the names of scheduled jobs with an empty env file.") readiness.add_argument("--queued", action="store_true", help="Only the names of the queued jobs, in the order they run.") + readiness.add_argument("--watched", action="store_true", + help="Name, account and watched directory of each job with a watch.") run_parser = sub.add_parser("run", help="Run a job as its timer does.") run_parser.add_argument("job") run_parser.add_argument("--dry-run", action="store_true", help="The job's own dry run; an alert is printed, not posted.") units_parser = sub.add_parser("units", help="Write each job's systemd drop-ins.") units_parser.add_argument("--out", required=True, metavar="DIR") + polkit_parser = sub.add_parser("polkit", help="Write the rule letting the Discord relay " + "start the jobs offered to /run.") + polkit_parser.add_argument("--out", required=True, metavar="FILE") + sub.add_parser("discord-register", help="Register the Discord slash commands, from " + "discord.env's variables.") argv = sys.argv[1:] if argv is None else list(argv) extra = argv[argv.index("--") + 1:] if "--" in argv else [] args = parser.parse_args(argv[:argv.index("--")] if "--" in argv else argv) @@ -250,6 +259,11 @@ def main(argv=None): if args.queued: print("\n".join(queue.jobs)) return + if args.watched: + for job in registry.load(): + if job.watch: + print(job.name, job.user, os.path.dirname(job.watch_glob)) + return for job in registry.load(): if args.ready or args.not_ready: if job.scheduled and ready(job) == args.ready: @@ -263,6 +277,16 @@ def main(argv=None): for path in units.write(registry.load(), registry.load_queue(), args.out): print(path) return + if args.command == "polkit": + from session_ops.ops import units + with open(args.out, "w", encoding="utf-8") as handle: + handle.write(units.polkit_rule(registry.load())) + print(args.out) + return + if args.command == "discord-register": + from session_ops.ops import discord_commands + discord_commands.main() + return try: job = registry.get(args.job) except KeyError: diff --git a/src/session_ops/ops/units.py b/src/session_ops/ops/units.py index 0d1ec64..d7b640b 100644 --- a/src/session_ops/ops/units.py +++ b/src/session_ops/ops/units.py @@ -4,9 +4,12 @@ drop-in carries only what jobs.toml says about it. Generated, so the registry is the one place a schedule, an account or a secret file is set. """ +import json import os HEADER = "# Generated by `session-ops units` from jobs.toml. Edit the registry, not this.\n" +# The account session-ops-discord.service runs as, and install.sh creates. +RELAY_USER = "opsbot" def service_dropin(job): @@ -25,6 +28,10 @@ def timer_dropin(schedule): return "\n".join([HEADER, "[Timer]", f"OnCalendar={schedule}"]) + "\n" +def path_dropin(pattern): + return "\n".join([HEADER, "[Path]", f"PathExistsGlob={pattern}"]) + "\n" + + def dropins(jobs, queue): """{relative path: content} for every job, and the queue's schedule. @@ -38,9 +45,26 @@ def dropins(jobs, queue): files[f"session-ops@{job.name}.service.d/job.conf"] = service_dropin(job) if job.schedule: files[f"session-ops@{job.name}.timer.d/schedule.conf"] = timer_dropin(job.schedule) + if job.watch: + files[f"session-ops@{job.name}.path.d/watch.conf"] = path_dropin(job.watch_glob) return files +def polkit_rule(jobs): + """Lets the Discord relay's account start the jobs offered to `/run`, and nothing else.""" + allowed = json.dumps([f"session-ops@{job.name}.service" for job in jobs if job.discord]) + return f"""// Generated by `session-ops polkit` from jobs.toml. Edit the registry, not this. +polkit.addRule(function(action, subject) {{ + if (action.id == "org.freedesktop.systemd1.manage-units" && + subject.user == "{RELAY_USER}" && + action.lookup("verb") == "start" && + {allowed}.indexOf(action.lookup("unit")) >= 0) {{ + return polkit.Result.YES; + }} +}}); +""" + + def write(jobs, queue, out): written = [] for relative, content in sorted(dropins(jobs, queue).items()): diff --git a/src/session_ops/platforms/mau.py b/src/session_ops/platforms/mau.py new file mode 100644 index 0000000..3562201 --- /dev/null +++ b/src/session_ops/platforms/mau.py @@ -0,0 +1,492 @@ +""" +Monthly active users, posted once a month: Android's and iOS's from the store exports +dropped into the job's inbox, iOS's scaled up by Apple's opt-in rate from the App Store +Connect API, the latest Android APKs' and Desktop release's downloads from GitHub, and +an F-Droid estimate. + + session-ops run mau [--dry-run] + /mau-upload file:.csv # in Discord, through session-ops-discord + rsync .csv root@:/var/lib/session-ops/mau/inbox/ + +Every export's daily figures merge into history.json, so an export may cover any range +and overlap the previous one. A month is posted once its last day is in the history for +every platform; until then, from the REMIND_DAY on, each run posts a reminder instead. +""" +import argparse +import csv +import glob +import gzip +import io +import json +import math +import os +import time +from datetime import date, datetime, timedelta +from zoneinfo import ZoneInfo + +import jwt + +from session_ops.ops.runner import step +from session_ops.shared import discord, http + +PLAY_COLUMN = ("Monthly Active Users (MAU) (Unique users, Per interval, Daily): " + "All countries / regions") +PLAY_DATE = "%b %d, %Y" +APPLE_COLUMN = "Active Last 30 Days" +APPLE_DATE = "%m/%d/%y" +PLATFORMS = ("android", "ios") +# Play's daily figures trail by about eight days, so a month's last day lands around the 9th. +REMIND_DAY = 10 +ZONE = ZoneInfo("Australia/Melbourne") +VERSION = 1 +HISTORY = "history.json" + +ANDROID_RELEASES = "https://api.github.com/repos/session-foundation/session-android/releases" +DESKTOP_RELEASES = "https://api.github.com/repos/session-foundation/session-desktop/releases" +FLATHUB = "https://flathub.org/api/v2/stats/network.loki.Session" +PLATFORM_EXTENSIONS = { + "linux": (".deb", ".AppImage", ".rpm", ".freebsd"), + "macos": (".dmg", ".zip"), + "windows": (".exe",), +} +# No store publishes F-Droid's downloads: this share of every other figure stands in for them. +FDROID_SHARE = 0.10 + +ASC_API = "https://api.appstoreconnect.apple.com" +IOS_APP_ID = "1470168868" +ASC_KEY_CREDENTIAL = "asc-key.p8" + + +class Rejected(ValueError): + pass + + +def read_rows(path): + try: + with open(path, encoding="utf-8-sig", newline="") as handle: + # strict: a quoted figure cut off at the end of the file is an error, not a smaller figure. + return list(csv.reader(handle, strict=True)) + except (UnicodeDecodeError, csv.Error) as exc: + raise Rejected(f"not a CSV export ({exc})") from None + except OSError as exc: + raise Rejected(f"unreadable ({exc.strerror})") from None + + +def parse_export(path): + """(platform, {ISO date: figure}) from a Play Console MAU export or an App Store + Connect Active in Last 30 Days export.""" + rows = read_rows(path) + first = rows[0] if rows and rows[0] else [""] + if first[0] == "Date" and PLAY_COLUMN in first: + return "android", parse_play(rows) + if first[0] == "Name": + return "ios", parse_apple(rows) + raise Rejected("neither a Play Console MAU export nor an App Store Connect Active in " + "Last 30 Days one") + + +def parse_play(rows): + column = rows[0].index(PLAY_COLUMN) + days = {} + for line, row in enumerate(rows[1:], start=2): + try: + day = datetime.strptime(row[0], PLAY_DATE).date().isoformat() + figure = row[column].replace(",", "") + except (IndexError, ValueError): + raise Rejected(f"line {line} is not a date and a figure: {row[:column + 1]}") from None + if not figure: + continue + if not figure.isdigit(): + raise Rejected(f"line {line} has {row[column]!r} for a figure") + days[day] = int(figure) + return non_empty(days) + + +def parse_apple(rows): + """The table below App Store Connect's Name, Start Date and End Date lines.""" + try: + header = next(i for i, row in enumerate(rows) if row[:1] == ["Date"]) + except StopIteration: + raise Rejected("an App Store Connect export with no Date column") from None + if rows[header] != ["Date", APPLE_COLUMN]: + raise Rejected(f"an App Store Connect export of {', '.join(rows[header][1:])}: " + f"export {APPLE_COLUMN} instead, daily") + days = {} + for line, row in enumerate(rows[header + 1:], start=header + 2): + try: + day = datetime.strptime(row[0], APPLE_DATE).date().isoformat() + value = float(row[1]) + except (IndexError, ValueError): + raise Rejected(f"line {line} is not a date and a figure: {row[:2]}") from None + if not math.isfinite(value) or value < 0: + raise Rejected(f"line {line} has {row[1]!r} for a figure") + # A day Apple is still counting comes as a fraction: kept, so it reads as provisional. + days[day] = int(value) if value.is_integer() else value + return non_empty(days) + + +def non_empty(days): + if not days: + raise Rejected("no figures in it") + return days + + +def read_inbox(inbox): + """(path, platform, days) for each export, oldest first, and (path, reason) for each + rejected one. + + glob skips dotfiles, so a file rsync or /mau-upload is still writing is never read. + """ + exports, rejected = [], [] + for path in sorted(glob.glob(os.path.join(inbox, "*.csv")), key=os.path.getmtime): + try: + exports.append((path, *parse_export(path))) + except Rejected as exc: + rejected.append((path, str(exc))) + return exports, rejected + + +def is_final(value): + return value is not None and float(value).is_integer() + + +def merge(history, exports): + """Merge each export into its platform's history, the later export winning, and record + in history["revisions"] each final figure a later export changed, until it is posted.""" + for _, platform, days in exports: + for day, figure in sorted(days.items()): + old = history[platform].get(day) + history[platform][day] = figure + if not is_final(old) or old == figure: + continue + key = f"{platform} {day}" + first = history["revisions"].get(key, [old])[0] + if first == figure: + history["revisions"].pop(key, None) + else: + history["revisions"][key] = [first, figure] + print(f"{key}: {old} revised to {figure}") + + +def shown_revisions(history, days): + """(platform, day, first, latest) for the recorded revisions of `days`.""" + revisions = [] + for key, (first, latest) in sorted(history["revisions"].items()): + platform, day = key.split(" ") + if day in days: + revisions.append((platform, day, first, latest)) + return revisions + + +def file_away(path, folder): + os.makedirs(folder, exist_ok=True) + stamp = time.strftime("%Y%m%dT%H%M%SZ", time.gmtime()) + os.replace(path, os.path.join(folder, f"{stamp}-{os.path.basename(path)}")) + + +def load_history(path): + """The history at `path`, empty if there is none yet. Unlike a digest's dedup cache it + cannot be rebuilt from a re-run, so an unreadable one stops the run rather than reset.""" + if not os.path.exists(path): + return {"version": VERSION, "android": {}, "ios": {}, "posted": [], "revisions": {}} + with open(path, encoding="utf-8") as handle: + data = json.load(handle) + if data.get("version") != VERSION: + raise RuntimeError(f"{path} is version {data.get('version')!r}, expected {VERSION}") + return data + + +def save_history(path, history): + temporary = f"{path}.tmp" + with open(temporary, "w", encoding="utf-8") as handle: + json.dump(history, handle, indent=2, sort_keys=True) + os.replace(temporary, path) + + +def previous_month_end(today): + return today.replace(day=1) - timedelta(days=1) + + +def figure(value): + return f"{round(value):,}" + + +def get_json(session, url): + resp = session.request("GET", url) + if resp.status_code != 200: + raise RuntimeError(f"HTTP {resp.status_code}: {resp.text[:200]}") + return resp.json() + + +def latest(releases): + return next(r for r in releases if not r["draft"] and not r["prerelease"]) + + +def release_day(release): + return release["published_at"].split("T")[0] + + +def flathub_installs_since(stats, day): + """Flathub installs from `day` on. Flathub builds from the GitHub .deb once, on its + own servers, so these are not already in the GitHub counts.""" + per_day = stats["installs_per_day"] + if day < min(per_day): + raise RuntimeError(f"Flathub keeps daily installs from {min(per_day)} only, " + f"after the release on {day}") + return sum(count for d, count in per_day.items() if d >= day) + + +def platform_totals(release, flathub_installs): + totals = {platform: sum(a["download_count"] for a in release["assets"] + if a["name"].endswith(extensions)) + for platform, extensions in PLATFORM_EXTENSIONS.items()} + totals["linux"] += flathub_installs + return totals + + +def desktop_downloads(session): + """(latest stable Desktop release, its downloads per platform).""" + release = latest(get_json(session, DESKTOP_RELEASES)) + flathub = flathub_installs_since(get_json(session, FLATHUB), release_day(release)) + return release, platform_totals(release, flathub) + + +def apk_downloads(session): + """(latest stable Android release, its APKs' downloads).""" + release = latest(get_json(session, ANDROID_RELEASES)) + return release, sum(a["download_count"] for a in release["assets"] + if a["name"].endswith(".apk")) + + +def asc_token(issuer, key_id, private_key): + now = int(time.time()) + return jwt.encode({"iss": issuer, "iat": now, "exp": now + 1200, "aud": "appstoreconnect-v1"}, + private_key, algorithm="ES256", headers={"kid": key_id, "typ": "JWT"}) + + +def asc_list(session, path): + url = ASC_API + path + while url: + page = get_json(session, url) + yield from page["data"] + url = page.get("links", {}).get("next") + + +def opt_in_days(asc, downloads, since): + """{ISO date: (first-time downloaders, those opting in)} from the App Opt In report of + every analytics request on the app, processed on or after `since`; where instances + overlap, the latest processed wins, since a later one carries late events. + + `downloads` fetches the segments: their URLs are pre-signed, so take no API token. + """ + requests = list(asc_list(asc, f"/v1/apps/{IOS_APP_ID}/analyticsReportRequests")) + stopped = [r["id"] for r in requests if r["attributes"]["accessType"] == "ONGOING" + and r["attributes"]["stoppedDueToInactivity"]] + if stopped: + raise RuntimeError(f"Apple stopped the ongoing analytics request {stopped[0]} for " + "inactivity; an Admin key must create a new one") + instances = [] + for request in requests: + for report in asc_list(asc, f"/v1/analyticsReportRequests/{request['id']}/reports" + "?filter[name]=App%20Opt%20In"): + instances += [i for i in asc_list(asc, f"/v1/analyticsReports/{report['id']}/instances" + "?filter[granularity]=DAILY&limit=200") + if i["attributes"]["processingDate"] >= since] + days = {} + for instance in sorted(instances, key=lambda i: i["attributes"]["processingDate"]): + for segment in asc_list(asc, f"/v1/analyticsReportInstances/{instance['id']}/segments"): + resp = downloads.request("GET", segment["attributes"]["url"]) + if resp.status_code != 200: + raise RuntimeError(f"HTTP {resp.status_code} for an App Opt In segment") + text = gzip.decompress(resp.content).decode("utf-8") + for row in csv.DictReader(io.StringIO(text), delimiter="\t"): + days[row["Date"]] = (int(row["Downloading Users"] or 0), + int(row["Users Opting-In"] or 0)) + return days + + +def opt_in(days, month_end): + """(first-time downloaders, those opting in) over the month, or None without any.""" + month = [v for d, v in days.items() if d[:7] == month_end.isoformat()[:7]] + downloading = sum(v[0] for v in month) + return (downloading, sum(v[1] for v in month)) if downloading else None + + +def ios_estimate(opted_in, opt_in_counts): + downloading, opting_in = opt_in_counts + return round(opted_in * downloading / opting_in) + + +LABELS = {"android": "Android", "ios": "iOS"} +HOW_TO_EXPORT = { + "android": "Play Console → Statistics → Saved reports → the MAU report, covering {day} → " + "Export report → CSV", + "ios": "App Store Connect → Analytics → Session → Metrics → Active in Last 30 Days, daily, " + "covering {day} → Export", +} + + +def with_change(current, before, month_end): + text = f"**{figure(current)}**" + if before: + change = current - before + text += (f" ({'+' if change >= 0 else '−'}{figure(abs(change))}, " + f"{change / before:+.1%} on {previous_month_end(month_end):%B})") + return text + + +def report_message(month_end, history, revisions, sources): + """`sources`: the desktop and APK (release, downloads) pairs, and Apple's opt-in + counts for this month and the one before, None where it has no first-time downloader.""" + day, before_day = month_end.isoformat(), previous_month_end(month_end).isoformat() + desktop_release, downloads = sources["desktop"] + apk_release, apks = sources["apks"] + counts, before_counts = sources["opt_in"] + if not counts or not counts[1]: + raise RuntimeError(f"Apple has no opt-in rate for {month_end:%B %Y}") + play = history["android"][day] + opted_in = history["ios"][day] + ios = ios_estimate(opted_in, counts) + ios_before = (ios_estimate(history["ios"][before_day], before_counts) + if before_counts and before_counts[1] and before_day in history["ios"] + else None) + desktop = downloads["linux"] + downloads["macos"] + downloads["windows"] + fdroid = round(FDROID_SHARE * (play + ios + apks + desktop)) + total = play + ios + apks + fdroid + desktop + on = f"{month_end:%-d %B}" + lines = [ + f"📊 **Monthly active users, {month_end:%B %Y}**", + f"Android, Play: {with_change(play, history['android'].get(before_day), month_end)}", + f"Android, outside Play: **{figure(apks + fdroid)}** " + f"(GitHub APKs {figure(apks)} · F-Droid {figure(fdroid)}, estimated)", + f"iOS: {with_change(ios, ios_before, month_end)}, estimated", + f"Desktop: **{figure(desktop)}** (Linux {figure(downloads['linux'])} · " + f"macOS {figure(downloads['macos'])} · Windows {figure(downloads['windows'])})", + f"**Total: {figure(total)}**", + f"-# Android, Play: Play Console MAU on {on}, users who opened Session in the 28 days " + "before.", + f"-# GitHub APKs: downloads of {apk_release['tag_name']}'s APKs since its release on " + f"{date.fromisoformat(release_day(apk_release)):%-d %B}, updates included.", + f"-# F-Droid publishes no counts: estimated as {FDROID_SHARE:.0%} of the other figures.", + f"-# iOS: Apple counts only devices that share analytics, {figure(opted_in)} active in " + f"the 30 days to {on}, scaled by {month_end:%B}'s opt-in rate from Apple's App Opt In " + f"report: {counts[1] / counts[0]:.1%}, {figure(counts[1])} of {figure(counts[0])} " + "first-time downloaders.", + f"-# Desktop: downloads of {desktop_release['tag_name']} since its release on " + f"{date.fromisoformat(release_day(desktop_release)):%-d %B}, updates included; " + "Desktop has no active-user count.", + ] + if revisions: + lines.append("-# The latest exports revised " + ", ".join( + f"{LABELS[p]} {date.fromisoformat(d):%-d %b} {figure(old)} → {figure(new)}" + for p, d, old, new in revisions)) + return "\n".join(lines) + + +def reminder_message(month_end, missing, provisional=()): + names = " and ".join(LABELS[p] for p in missing) + day = f"{month_end:%-d %B}" + return "\n".join([ + f"⏰ **{names} active users for {month_end:%B %Y} {'are' if len(missing) > 1 else 'is'} " + "missing.** Export:", + *(f"- {LABELS[p]}: " + HOW_TO_EXPORT[p].format(day=day) + + (" (the last export had it still being counted)" if p in provisional else "") + for p in missing), + "Then hand each to `/mau-upload` here, and the figures post as soon as the last one " + "lands.", + ]) + + +def read_sources(args, month_end): + step("reading GitHub and Flathub downloads") + github = http.Session() + github.headers.update({"Accept": "application/vnd.github.v3+json"}) + sources = {"desktop": desktop_downloads(github), "apks": apk_downloads(github)} + + step("reading Apple's opt-in rate") + with open(args.asc_key, encoding="utf-8") as handle: + key = handle.read() + if not key.strip(): + raise SystemExit(f"{args.asc_key} is empty: put the App Store Connect key's .p8 at " + f"/etc/session-ops/{ASC_KEY_CREDENTIAL}, mode 600") + token = asc_token(os.environ["ASC_ISSUER_ID"], os.environ["ASC_KEY_ID"], key) + asc = http.Session() + asc.headers.update({"Authorization": f"Bearer {token}"}) + month_before = previous_month_end(month_end) + days = opt_in_days(asc, http.Session(), since=month_before.replace(day=1).isoformat()) + sources["opt_in"] = (opt_in(days, month_end), opt_in(days, month_before)) + return sources + + +def main(argv=None): + parser = argparse.ArgumentParser(description=__doc__.strip().split("\n")[0]) + parser.add_argument("--state", required=True, metavar="DIR", + help="Holds inbox/, done/, rejected/ and history.json.") + parser.add_argument("--dry-run", action="store_true", + help="Print what would be posted; move and write nothing.") + parser.add_argument("--asc-key", metavar="PATH", + default=os.path.join(os.environ.get("CREDENTIALS_DIRECTORY", ""), + ASC_KEY_CREDENTIAL), + help="The App Store Connect API key, for Apple's opt-in rate.") + args = parser.parse_args(argv) + + inbox = os.path.join(args.state, "inbox") + history_path = os.path.join(args.state, HISTORY) + step("reading the inbox") + try: + history = load_history(history_path) + exports, rejected = read_inbox(inbox) + merge(history, exports) + if not args.dry_run: + # Saved before the files move: a run stopped in between reads them again, harmlessly. + save_history(history_path, history) + for path, _, _ in exports: + file_away(path, os.path.join(args.state, "done")) + for path, _ in rejected: + file_away(path, os.path.join(args.state, "rejected")) + except Exception: + # The path unit restarts the job while a file matches its glob, then stops watching. + if not args.dry_run: + for path in glob.glob(os.path.join(inbox, "*.csv")): + file_away(path, os.path.join(args.state, "rejected")) + print("The inbox's exports are in rejected/: move them back once this is fixed.") + raise + + today = datetime.now(ZONE).date() + month_end = previous_month_end(today) + month = month_end.strftime("%Y-%m") + missing = [p for p in PLATFORMS if not is_final(history[p].get(month_end.isoformat()))] + provisional = [p for p in missing if month_end.isoformat() in history[p]] + message = None + if month in history["posted"]: + print(f"{month} already posted.") + elif not missing: + revisions = shown_revisions(history, {month_end.isoformat(), + previous_month_end(month_end).isoformat()}) + message = report_message(month_end, history, revisions, read_sources(args, month_end)) + elif today.day >= REMIND_DAY: + message = reminder_message(month_end, missing, provisional) + else: + print(f"Waiting for {month_end} from {', '.join(missing)}; " + f"reminders start on the {REMIND_DAY}th.") + + if message: + step("posting to Discord") + if args.dry_run: + print(message) + else: + payload = {"content": message, "allowed_mentions": {"parse": []}} + if not discord.post_to_discord(http.Session(), os.environ["MAU_DISCORD_WEBHOOK_URL"], + [payload]): + raise RuntimeError("Discord did not accept the message") + if not missing: + history["posted"].append(month) + history["revisions"] = {} + save_history(history_path, history) + if rejected: + raise SystemExit("rejected " + "; ".join( + f"{os.path.basename(path)}: {reason}" for path, reason in rejected)) + + +if __name__ == "__main__": + main() diff --git a/src/session_ops/platforms/release_stats.py b/src/session_ops/platforms/release_stats.py deleted file mode 100644 index 844de28..0000000 --- a/src/session_ops/platforms/release_stats.py +++ /dev/null @@ -1,109 +0,0 @@ -""" -Download counts of the last ten Desktop and Android releases, per installer, as two -CSV files. On demand only: nothing schedules it. - - session-ops run release-stats # CSVs kept under the job's runs/ - uv run python -m session_ops.platforms.release_stats --out . -""" -import argparse -import os -import time -from datetime import datetime, timezone - -from session_ops.ops.runner import step -from session_ops.shared import http - -API = "https://api.github.com/repos/session-foundation/{repo}/releases" -RELEASES = 10 - -DESKTOP_HEADER = ("version,snapshot_date,release_date,.deb,.appimage,.rpm,.dmg_arm64," - ".dmg_x64,.zip_arm64,.zip_x64,.exe") -ANDROID_HEADER = ("version,snapshot_date,release_date,.aab,.apk_arm64,.apk_armv7a," - ".apk_universal_huawei,.apk_universal_play,.apk_x86,.apk_x86_64") - - -def downloads(assets, keep): - return sum(a["download_count"] for a in assets if keep(a["name"])) - - -def desktop_row(release): - assets = release["assets"] - return [ - downloads(assets, lambda n: n.endswith(".deb")), - downloads(assets, lambda n: n.endswith(".AppImage")), - downloads(assets, lambda n: n.endswith(".rpm")), - downloads(assets, lambda n: n.endswith(".dmg") and "arm64" in n), - downloads(assets, lambda n: n.endswith(".dmg") and "x64" in n), - downloads(assets, lambda n: n.endswith(".zip") and "arm64" in n), - downloads(assets, lambda n: n.endswith(".zip") and "x64" in n), - downloads(assets, lambda n: n.endswith(".exe")), - ] - - -def android_row(release): - assets = release["assets"] - - def apk(arch): - return downloads(assets, lambda n: n.endswith(".apk") and arch in n - and "play-release" in n) - - return [ - downloads(assets, lambda n: n.endswith(".aab") and "play-release" in n), - apk("arm64-v8a"), - apk("armeabi-v7a"), - downloads(assets, lambda n: n.endswith(".apk") and "universal" in n - and "huawei-release" in n), - downloads(assets, lambda n: n.endswith(".apk") and "universal" in n - and "play-release" in n), - # x86 without x86_64, which contains it. - downloads(assets, lambda n: n.endswith(".apk") and "x86" in n - and "x86_64" not in n and "play-release" in n), - apk("x86_64"), - ] - - -def csv(releases, header, row, snapshot): - lines = [header] - for release in releases[:RELEASES]: - version = release["tag_name"][1:] if release["tag_name"].startswith("v") \ - else release["tag_name"] - lines.append(",".join(str(v) for v in [version, snapshot, - release["published_at"].split("T")[0], - *row(release)])) - return "\n".join(lines) - - -def fetch(session, repo): - resp = session.request("GET", API.format(repo=repo)) - if resp.status_code != 200: - raise RuntimeError(f"HTTP {resp.status_code}: {resp.text[:200]}") - return resp.json() - - -def main(argv=None): - parser = argparse.ArgumentParser(description=__doc__.strip().split("\n")[0]) - parser.add_argument("--out", default=".", metavar="DIR", - help="Where the CSVs go; a directory per run is made under it.") - parser.add_argument("--dry-run", action="store_true", help="Print the CSVs only.") - args = parser.parse_args(argv) - - session = http.Session() - session.headers.update({"Accept": "application/vnd.github.v3+json"}) - snapshot = datetime.now(timezone.utc).date().isoformat() - out = os.path.join(args.out, time.strftime("%Y%m%dT%H%M%SZ", time.gmtime())) - for repo, header, row in (("session-desktop", DESKTOP_HEADER, desktop_row), - ("session-android", ANDROID_HEADER, android_row)): - step(repo) - content = csv(fetch(session, repo), header, row, snapshot) - print(f"# {repo}\n{content}\n") - if not args.dry_run: - os.makedirs(out, exist_ok=True) - with open(os.path.join(out, f"{repo}-release-stats.csv"), "w", - encoding="utf-8") as handle: - handle.write(content) - if not args.dry_run: - print(f"Written to {out}") - - -if __name__ == "__main__": - main() diff --git a/tests/goldens/units/dropins.txt b/tests/goldens/units/dropins.txt index 7f3cb1a..f4e6946 100644 --- a/tests/goldens/units/dropins.txt +++ b/tests/goldens/units/dropins.txt @@ -44,17 +44,30 @@ Group=ghdigest EnvironmentFile=/etc/session-ops/github-prs.env TimeoutStartSec=15min -==> session-ops@release-stats.service.d/job.conf <== +==> session-ops@mau.path.d/watch.conf <== +# Generated by `session-ops units` from jobs.toml. Edit the registry, not this. + +[Path] +PathExistsGlob=/var/lib/session-ops/mau/inbox/*.csv + +==> session-ops@mau.service.d/job.conf <== # Generated by `session-ops units` from jobs.toml. Edit the registry, not this. [Unit] -Description=Download counts of the latest Desktop and Android releases +Description=Monthly active users per platform, posted once a month [Service] User=sessionops Group=sessionops -EnvironmentFile=/etc/session-ops/alerts.env +EnvironmentFile=/etc/session-ops/mau.env TimeoutStartSec=5min +LoadCredential=asc-key.p8:/etc/session-ops/asc-key.p8 + +==> session-ops@mau.timer.d/schedule.conf <== +# Generated by `session-ops units` from jobs.toml. Edit the registry, not this. + +[Timer] +OnCalendar=*-*-10 11:00 Australia/Melbourne ==> session-ops@session-ops-silence.service.d/job.conf <== # Generated by `session-ops units` from jobs.toml. Edit the registry, not this. diff --git a/tests/goldens/units/polkit.rules b/tests/goldens/units/polkit.rules new file mode 100644 index 0000000..a9a16f1 --- /dev/null +++ b/tests/goldens/units/polkit.rules @@ -0,0 +1,9 @@ +// Generated by `session-ops polkit` from jobs.toml. Edit the registry, not this. +polkit.addRule(function(action, subject) { + if (action.id == "org.freedesktop.systemd1.manage-units" && + subject.user == "opsbot" && + action.lookup("verb") == "start" && + ["session-ops@crowdin-duplicates.service", "session-ops@crowdin-sync.service"].indexOf(action.lookup("unit")) >= 0) { + return polkit.Result.YES; + } +}); diff --git a/tests/ops/test_deploy.py b/tests/ops/test_deploy.py index ffe4a0f..85625ad 100644 --- a/tests/ops/test_deploy.py +++ b/tests/ops/test_deploy.py @@ -14,7 +14,7 @@ import tomllib import unittest -from session_ops.ops import registry, units +from session_ops.ops import discord_commands, registry, units from tests.golden import assert_golden ROOT = os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) @@ -62,6 +62,16 @@ def test_a_scheduled_job_gets_a_timer_and_an_unscheduled_one_does_not(self): self.assertEqual(f"session-ops@{job.name}.timer.d/schedule.conf" in files, bool(job.schedule)) + def test_a_watched_job_gets_a_path_unit_and_moves_its_files_out_of_the_watch(self): + files = units.dropins(registry.load(), registry.load_queue()) + for job in registry.load(): + with self.subTest(job=job.name): + watch = files.get(f"session-ops@{job.name}.path.d/watch.conf") + self.assertEqual(watch is not None, bool(job.watch)) + if job.watch: + self.assertIn(f"\nPathExistsGlob={job.watch_glob}\n", watch) + self.assertIn("Unit=session-ops@%i.service\n", unit_text("session-ops@.path")) + def test_each_queued_job_runs_after_every_one_before_it(self): queue = registry.load_queue() files = units.dropins(registry.load(), queue) @@ -79,6 +89,23 @@ def test_the_generated_dropins(self): text = "".join(f"==> {path} <==\n{files[path]}\n" for path in sorted(files)) assert_golden(self, "units/dropins.txt", text) + def test_the_generated_polkit_rule(self): + assert_golden(self, "units/polkit.rules", units.polkit_rule(registry.load())) + + def test_the_polkit_rule_names_the_account_the_discord_relay_runs_as(self): + self.assertIn(f"\nUser={units.RELAY_USER}\n", unit_text("session-ops-discord.service")) + self.assertIn(f'subject.user == "{units.RELAY_USER}"', units.polkit_rule([])) + + def test_the_discord_relay_can_write_mau_s_inbox_and_nothing_else(self): + inbox = os.path.dirname(registry.get(discord_commands.MAU_JOB).watch_glob) + writable = re.findall(r"^ReadWritePaths=-?(.*)$", + unit_text("session-ops-discord.service"), re.MULTILINE) + self.assertEqual(writable, [inbox]) + with open(os.path.join(ROOT, "deploy", "install.sh"), encoding="utf-8") as handle: + self.assertIn(f'if [ "$job" = {discord_commands.MAU_JOB} ]; then\n' + f' install -d -o "$user" -g {units.RELAY_USER} -m 770', + handle.read()) + class TestInstallScript(unittest.TestCase): @unittest.skipIf(ROOT == "/opt/session-ops", "this checkout is the one it installs from") diff --git a/tests/ops/test_discord_relay.py b/tests/ops/test_discord_relay.py new file mode 100644 index 0000000..e11642f --- /dev/null +++ b/tests/ops/test_discord_relay.py @@ -0,0 +1,297 @@ +""" + uv run python -m unittest tests.ops.test_discord_relay + +Offline: systemctl, Discord's CDN and its webhook API are fakes, and requests go +through FastAPI's TestClient. The relay can start jobs on the host, so most of this is +about who it refuses. +""" +import json +import os +import shutil +import subprocess +import tempfile +import time +import unittest +from unittest import mock + +from fastapi.testclient import TestClient +from nacl.signing import SigningKey + +from session_ops.ops import discord_commands, discord_relay as relay, registry +from session_ops.platforms import mau +from session_ops.shared.testing import FakeResponse, FakeSession, Patched + +KEY = SigningKey.generate() +GUILD = "111" +ALLOWED_USER = "222" +ALLOWED_ROLE = "333" +ENV = {"DISCORD_PUBLIC_KEY": KEY.verify_key.encode().hex(), "DISCORD_GUILD_ID": GUILD, + "ALLOWED_USER_IDS": ALLOWED_USER, "ALLOWED_ROLE_IDS": ALLOWED_ROLE} +EXPORT = f'Date,"{mau.PLAY_COLUMN}"\n"Sep 29, 2026","1,200"\n"Sep 30, 2026","1,250"\n' +CDN_URL = "https://cdn.discordapp.com/attachments/1/2/All%20countries.csv?ex=1" + + +class FakeSystemd: + def __init__(self, running=(), waiting=(), start_error=None, queue_at=None): + self.running, self.waiting, self.start_error = set(running), set(waiting), start_error + self.queue_at = queue_at + self.started, self.locked = [], [] + + def __call__(self, *args): + out, code, err = "", 0, "" + if args[0] == "list-jobs": + out = "".join(f"7 {unit} start waiting\n" for unit in self.waiting) + elif args[0] == "is-active": + out = "".join(("active" if u in self.running else "inactive") + "\n" + for u in args[1:]) + elif args[0] == "list-timers": + out = json.dumps([] if self.queue_at is None else [ + {"next": int(self.queue_at * 1e6), "unit": args[-1]}]) + elif args[0] == "start": + self.started.append(args[-1]) + self.locked.append(relay._starting.locked()) + code, err = (1, self.start_error) if self.start_error else (0, "") + return subprocess.CompletedProcess(["systemctl", *args], code, out, err) + + +def command(name, options=(), *, user=ALLOWED_USER, roles=(), guild=GUILD, resolved=None): + data = {"name": name, "options": [{"name": k, "value": v} for k, v in options]} + if resolved: + data["resolved"] = resolved + return {"type": relay.INTERACTION_COMMAND, "id": "999", "application_id": "app", + "token": "tok", "guild_id": guild, "data": data, + "member": {"user": {"id": user}, "roles": list(roles)}} + + +def upload(attachment): + return command(discord_commands.MAU_UPLOAD, [("file", "55")], + resolved={"attachments": {"55": attachment}}) + + +def attachment(**overrides): + return {"id": "55", "filename": "All countries _ regions.csv", "size": len(EXPORT), + "url": CDN_URL, **overrides} + + +def post(body, *, key=KEY, age=0, tamper=False): + raw = json.dumps(body).encode() + timestamp = str(int(time.time() - age)) + signature = key.sign(timestamp.encode() + raw).signature.hex() + if tamper: + raw = raw.replace(b"}", b" }", 1) + return TestClient(relay.app).post("/discord/interactions", content=raw, headers={ + "x-signature-ed25519": signature, "x-signature-timestamp": timestamp}) + + +class RelayCase(unittest.TestCase): + def setUp(self): + patcher = mock.patch.dict(os.environ, ENV) + patcher.start() + self.addCleanup(patcher.stop) + self.systemd = FakeSystemd() + patched = Patched(relay, systemctl=self.systemd) + patched.__enter__() + self.addCleanup(patched.__exit__) + + def content(self, response): + self.assertEqual(response.status_code, 200) + return response.json()["data"]["content"] + + +class TestGate(RelayCase): + def test_a_ping_is_answered(self): + self.assertEqual(post({"type": relay.INTERACTION_PING}).json(), {"type": relay.RESPONSE_PONG}) + + def test_an_unsigned_tampered_stale_or_foreign_request_is_refused(self): + for case in ({"tamper": True}, {"age": relay.MAX_SIGNATURE_AGE_SECONDS + 5}, + {"age": -relay.MAX_CLOCK_SKEW_SECONDS - 5}, {"key": SigningKey.generate()}): + with self.subTest(**{k: str(v) for k, v in case.items()}): + self.assertEqual(post({"type": relay.INTERACTION_PING}, **case).status_code, 401) + + def test_no_public_key_refuses_everything(self): + with mock.patch.dict(os.environ, {"DISCORD_PUBLIC_KEY": ""}): + self.assertEqual(post({"type": relay.INTERACTION_PING}).status_code, 401) + + def test_another_server_or_a_dm_is_refused(self): + for guild in ("444", None): + with self.subTest(guild=guild): + text = self.content(post(command("run", [("job", "crowdin-sync")], guild=guild))) + self.assertIn("do not work here", text) + self.assertEqual(self.systemd.started, []) + + def test_someone_on_neither_list_is_refused(self): + text = self.content(post(command("run", [("job", "crowdin-sync")], user="5", roles=["6"]))) + self.assertIn("not on the list", text) + self.assertEqual(self.systemd.started, []) + + def test_empty_allowlists_refuse_everybody(self): + with mock.patch.dict(os.environ, {"ALLOWED_USER_IDS": "", "ALLOWED_ROLE_IDS": ""}): + text = self.content(post(command("run", [("job", "crowdin-sync")]))) + self.assertIn("Neither", text) + + def test_a_role_on_the_list_is_enough(self): + post(command("run", [("job", "crowdin-sync")], user="5", roles=[ALLOWED_ROLE])) + self.assertEqual(self.systemd.started, ["session-ops@crowdin-sync.service"]) + + +class TestRun(RelayCase): + def test_a_discord_job_is_started_and_the_channel_told_who_started_it(self): + response = post(command("run", [("job", "crowdin-duplicates")])) + self.assertEqual(self.systemd.started, ["session-ops@crowdin-duplicates.service"]) + data = response.json()["data"] + self.assertIn(f"<@{ALLOWED_USER}> started **crowdin-duplicates**", data["content"]) + self.assertNotIn("flags", data) + self.assertEqual(data["allowed_mentions"], {"parse": []}) + + def test_a_job_not_marked_discord_is_refused(self): + for name in ("token-expiry", "nope"): + with self.subTest(job=name): + self.assertIn("not a job /run can start", + self.content(post(command("run", [("job", name)])))) + self.assertEqual(self.systemd.started, []) + + def test_a_running_or_waiting_queue_job_holds_it_off(self): + queued = registry.load_queue().jobs + for systemd in (FakeSystemd(running=[f"session-ops@{queued[0]}.service"]), + FakeSystemd(waiting=[f"session-ops@{queued[-1]}.service"]), + FakeSystemd(running=["session-ops@crowdin-sync.service"])): + with self.subTest(running=systemd.running, waiting=systemd.waiting), \ + Patched(relay, systemctl=systemd): + self.assertIn("running or waiting", self.content( + post(command("run", [("job", "crowdin-sync")])))) + self.assertEqual(systemd.started, []) + + def test_the_check_and_the_start_happen_under_one_lock(self): + post(command("run", [("job", "crowdin-sync")])) + self.assertEqual(self.systemd.locked, [True]) + self.assertFalse(relay._starting.locked()) + + def test_a_run_that_could_still_be_going_when_the_queue_starts_is_refused(self): + timeout = registry.get("crowdin-sync").timeout_seconds + systemd = FakeSystemd(queue_at=time.time() + timeout - 60) + with Patched(relay, systemctl=systemd): + text = self.content(post(command("run", [("job", "crowdin-sync")]))) + self.assertIn(f"The queue runs **crowdin-sync** at ", text) + self.assertEqual(systemd.started, []) + + def test_a_run_sure_to_end_before_the_queue_starts_goes_ahead(self): + timeout = registry.get("crowdin-duplicates").timeout_seconds + systemd = FakeSystemd(queue_at=time.time() + timeout + 60) + with Patched(relay, systemctl=systemd): + post(command("run", [("job", "crowdin-duplicates")])) + self.assertEqual(systemd.started, ["session-ops@crowdin-duplicates.service"]) + + def test_a_disabled_queue_timer_holds_nothing_off(self): + post(command("run", [("job", "crowdin-sync")])) + self.assertEqual(self.systemd.started, ["session-ops@crowdin-sync.service"]) + + def test_a_refusal_from_systemd_is_passed_on(self): + systemd = FakeSystemd(start_error="Failed to start: Access denied") + with Patched(relay, systemctl=systemd): + text = self.content(post(command("run", [("job", "crowdin-sync")]))) + self.assertIn("Access denied", text) + + def test_systemctl_missing_is_reported_not_raised(self): + def missing(*args): + raise FileNotFoundError("systemctl") + with Patched(relay, systemctl=missing): + self.assertIn("Could not ask systemd", + self.content(post(command("run", [("job", "crowdin-sync")])))) + + +class TestMauUpload(RelayCase): + def setUp(self): + super().setUp() + self.inbox = tempfile.mkdtemp() + self.addCleanup(shutil.rmtree, self.inbox, True) + self.cdn = FakeResponse({}) + self.cdn.text = EXPORT + self.session = FakeSession([self.cdn, FakeResponse({})]) + patched = Patched(relay, inbox=lambda: self.inbox) + patched.__enter__() + self.addCleanup(patched.__exit__) + http = Patched(relay.http, Session=lambda **kwargs: self.session) + http.__enter__() + self.addCleanup(http.__exit__) + + def told(self): + method, url, kwargs = self.session.calls[-1] + self.assertEqual((method, url), ("PATCH", f"{discord_commands.API}/webhooks/app/tok" + "/messages/@original")) + return kwargs["json"]["content"] + + def test_an_export_lands_in_the_inbox_under_a_name_mau_reads(self): + response = post(upload(attachment())) + self.assertEqual(response.json(), {"type": relay.RESPONSE_DEFERRED}) + self.assertEqual(os.listdir(self.inbox), ["discord-999.csv"]) + path = os.path.join(self.inbox, "discord-999.csv") + self.assertEqual(mau.parse_export(path), + ("android", {"2026-09-29": 1200, "2026-09-30": 1250})) + self.assertEqual(os.stat(path).st_mode & 0o777, 0o644) + self.assertIn("Android export, 2 days from 2026-09-29 to 2026-09-30", self.told()) + + def test_an_app_store_connect_export_is_named_ios(self): + self.cdn.text = ('Name,"Session - Private Messenger"\nStart Date,9/1/26\n' + 'End Date,9/30/26\n\nDate,Active Last 30 Days\n9/30/26,42000.0\n') + post(upload(attachment())) + self.assertIn("iOS export, 1 day from 2026-09-30", self.told()) + + def test_a_file_mau_would_reject_is_refused_and_leaves_nothing_behind(self): + self.cdn.text = 'Date,Installs\n"Sep 30, 2026",4\n' + post(upload(attachment())) + self.assertEqual(os.listdir(self.inbox), []) + self.assertIn("not the export mau reads", self.told()) + + def test_the_cdn_failing_is_reported(self): + self.cdn.status_code = 404 + post(upload(attachment())) + self.assertEqual(os.listdir(self.inbox), []) + self.assertIn("answered 404", self.told()) + + def test_what_cannot_be_an_export_is_refused_before_downloading(self): + for bad, why in ((attachment(filename="mau.xlsx"), "export a .csv"), + (attachment(size=relay.MAX_UPLOAD_BYTES + 1), "larger"), + (attachment(url="https://evil.test/a.csv"), "not on Discord's CDN"), + (attachment(url=CDN_URL.replace("https", "http")), "not on Discord's CDN"), + (None, "no file")): + with self.subTest(why=why): + self.assertIn(why, self.content(post(upload(bad)))) + self.assertEqual(self.session.calls, []) + self.assertEqual(os.listdir(self.inbox), []) + + def test_someone_not_allowed_cannot_upload(self): + self.assertIn("not on the list", self.content(post( + {**upload(attachment()), "member": {"user": {"id": "5"}, "roles": []}}))) + self.assertEqual(self.session.calls, []) + + +class TestCommands(unittest.TestCase): + def test_run_offers_exactly_the_discord_jobs_and_both_start_hidden(self): + jobs = registry.load() + run, upload_command = discord_commands.commands(jobs) + offered = [choice["value"] for choice in run["options"][0]["choices"]] + self.assertEqual(offered, [job.name for job in jobs if job.discord]) + self.assertIn("crowdin-sync", offered) + for item in (run, upload_command): + self.assertEqual(item["default_member_permissions"], "0") + for choice in run["options"][0]["choices"]: + self.assertLessEqual(len(choice["name"]), discord_commands.MAX_CHOICE_NAME) + + def test_registration_replaces_the_guild_commands(self): + session = FakeSession([FakeResponse([{"name": "run"}, {"name": "mau-upload"}])]) + names = discord_commands.register(session, "app", GUILD, "bot", registry.load()) + method, url, kwargs = session.calls[0] + self.assertEqual((method, url), ("PUT", f"{discord_commands.API}/applications/app" + f"/guilds/{GUILD}/commands")) + self.assertEqual(kwargs["headers"], {"Authorization": "Bot bot"}) + self.assertEqual(names, ["run", "mau-upload"]) + + def test_a_refused_registration_exits_saying_why(self): + session = FakeSession([FakeResponse({"message": "Missing Access"}, status_code=403)]) + with self.assertRaises(SystemExit) as raised: + discord_commands.register(session, "app", GUILD, "bot", registry.load()) + self.assertIn("Missing Access", str(raised.exception)) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/ops/test_registry.py b/tests/ops/test_registry.py index 88842f3..cd1a94e 100644 --- a/tests/ops/test_registry.py +++ b/tests/ops/test_registry.py @@ -40,6 +40,16 @@ def test_state_is_substituted_and_dry_run_appended(self): self.assertEqual(job.argv(dry_run=True, state_dir="/x"), ["--state", "/x/s.json", "--dry-run"]) + def test_a_watch_is_a_glob_under_the_state_directory(self): + job = load(VALID + 'watch = "inbox/*.csv"\n')[0] + self.assertEqual(job.watch_glob, "/var/lib/session-ops/a/inbox/*.csv") + self.assertIsNone(load(VALID)[0].watch_glob) + + def test_a_watch_outside_the_state_directory_is_refused(self): + for watch in ("/tmp/*.csv", "../b/*.csv"): + with self.subTest(watch=watch), self.assertRaises(ValueError): + load(VALID + f'watch = "{watch}"\n') + def test_a_missing_field_is_refused(self): with self.assertRaises(ValueError): load(VALID.replace('user = "u"\n', "")) @@ -70,6 +80,30 @@ def test_a_name_listed_twice_is_refused(self): with self.assertRaises(ValueError): load(VALID + VALID) + def test_a_job_is_offered_to_discord_only_when_it_says_so(self): + self.assertFalse(load(VALID)[0].discord) + self.assertTrue(load(VALID + "discord = true\n")[0].discord) + + def test_a_discord_flag_that_is_not_a_boolean_is_refused(self): + with self.assertRaises(ValueError): + load(VALID + 'discord = "yes"\n') + + def test_more_jobs_than_discord_offers_choices_for_is_refused(self): + many = "".join(VALID.replace('"a"', f'"j{i}"') + "discord = true\n" + for i in range(registry.MAX_DISCORD_JOBS + 1)) + with self.assertRaises(ValueError): + load(many) + + def test_a_timeout_is_read_as_a_systemd_time_span(self): + self.assertEqual(load(VALID)[0].timeout_seconds, 1800) + self.assertEqual(load(VALID + 'timeout = "1h 30min"\n')[0].timeout_seconds, 5400) + self.assertEqual(registry.span_seconds("90"), 90) + + def test_a_timeout_systemd_would_not_read_is_refused(self): + for timeout in ("soon", "2x", "1h later", ""): + with self.subTest(timeout=timeout), self.assertRaises(ValueError): + load(VALID + f'timeout = "{timeout}"\n') + def test_the_shipped_registry_loads(self): names = [job.name for job in registry.load()] self.assertIn("zendesk-digest", names) diff --git a/tests/test_mau.py b/tests/test_mau.py new file mode 100644 index 0000000..41aa9ea --- /dev/null +++ b/tests/test_mau.py @@ -0,0 +1,472 @@ +""" + uv run python -m unittest tests.test_mau +""" +import contextlib +import gzip +import io +import json +import os +import tempfile +import unittest +from datetime import date, datetime +from unittest import mock + +from session_ops.platforms import mau +from session_ops.shared.testing import FakeResponse, FakeSession + + +class GzipResponse(FakeResponse): + def __init__(self, text): + super().__init__({}) + self._content = gzip.compress(text.encode()) + + @property + def content(self): + return self._content + +PLAY_HEADER = f'Date,"{mau.PLAY_COLUMN}",Notes\n' + + +def asset(name, count): + return {"name": name, "download_count": count} + + +def release(tag, *assets, published="2026-07-10", prerelease=False): + return {"tag_name": tag, "published_at": f"{published}T00:00:00Z", "draft": False, + "prerelease": prerelease, "assets": list(assets)} + + +LATEST = release( + "v1.18.1", + asset("session-desktop-linux-amd64-1.18.1.deb", 5), + asset("session-desktop-linux-x86_64-1.18.1.AppImage", 7), + asset("session-desktop-linux-x86_64-1.18.1.rpm", 1), + asset("session-desktop-linux-x64-1.18.1.freebsd", 2), + asset("session-desktop-mac-arm64-1.18.1.dmg", 30), + asset("session-desktop-mac-arm64-1.18.1.dmg.blockmap", 900), + asset("session-desktop-mac-x64-1.18.1.zip", 4), + asset("session-desktop-win-x64-1.18.1.exe", 100), + asset("session-desktop-win-x64-1.18.1.exe.blockmap", 900), + asset("latest.yml", 900), + asset("latest-linux.yml", 900), + asset("signature.asc", 900)) + +DESKTOP = [ + release("v1.19.0", asset("session-desktop-win-x64-1.19.0.exe", 9), + published="2026-10-01", prerelease=True), + LATEST, + release("v1.18.0", asset("session-desktop-win-x64-1.18.0.exe", 50), + published="2026-04-09"), +] + +ANDROID = [release("1.33.6", asset("session-1.33.6-arm64-v8a-play-release.apk", 9), + prerelease=True), + release("1.33.5", asset("app-play-release.aab", 500), + asset("session-1.33.5-arm64-v8a-play-release.apk", 60), + asset("session-1.33.5-universal-huawei-release.apk", 4), + asset("signature.asc", 500), published="2026-07-13")] + +FLATHUB = {"installs_per_day": {"2026-07-11": 10, "2026-07-09": 1000, "2026-07-10": 5}} + + +def play(*rows): + return PLAY_HEADER + "".join(f'"{day}","{value}",{note}\n' for day, value, note in rows) + + +def apple(*rows, column=mau.APPLE_COLUMN): + return ('Name,"Session - Private Messenger"\nStart Date,8/31/26\nEnd Date,9/30/26\n\n' + f"Date,{column}\n" + "".join(f"{day},{value}\n" for day, value in rows)) + + +# LATEST's downloads with FLATHUB's 15 installs since its release: 30 + 34 + 100. +DOWNLOADS = (LATEST, {"linux": 30, "macos": 34, "windows": 100}) +ANDROID_SEPTEMBER = play(("Aug 31, 2026", "100,000", ""), + ("Sep 29, 2026", "104,500", "Rollout of release: 1.32.1 at 5%."), + ("Sep 30, 2026", "105,000", "")) +IOS_SEPTEMBER = apple(("8/31/26", "40000.0"), ("9/29/26", "41500.0"), ("9/30/26", "42000.0")) +HISTORY = {"android": {"2026-08-31": 100000, "2026-09-30": 105000}, + "ios": {"2026-08-31": 40000, "2026-09-30": 42000}} +# iOS: 42,000 opted in at 1 in 4 is 168,000; August's 40,000 at 1 in 5 is 200,000. +# F-Droid: 10% of 105,000 + 168,000 + 64 + 164. +SOURCES = {"desktop": DOWNLOADS, "apks": (ANDROID[1], 64), + "opt_in": ((400, 100), (500, 100))} + + +class ParseTest(unittest.TestCase): + def parse(self, text): + with tempfile.NamedTemporaryFile("w", suffix=".csv", delete=False) as handle: + handle.write(text) + self.addCleanup(os.remove, handle.name) + return mau.parse_export(handle.name) + + def test_a_play_export_is_android_without_thousands_separators(self): + self.assertEqual(self.parse(ANDROID_SEPTEMBER), + ("android", {"2026-08-31": 100000, "2026-09-29": 104500, + "2026-09-30": 105000})) + + def test_an_app_store_connect_export_is_ios_past_its_preamble(self): + self.assertEqual(self.parse(IOS_SEPTEMBER), + ("ios", {"2026-08-31": 40000, "2026-09-29": 41500, "2026-09-30": 42000})) + + def test_a_day_apple_is_still_counting_keeps_its_fraction(self): + self.assertEqual(self.parse(apple(("10/6/26", "19000.0"), ("10/7/26", "18973.094")))[1], + {"2026-10-06": 19000, "2026-10-07": 18973.094}) + + def test_an_unreadable_file_is_a_rejection(self): + with self.assertRaisesRegex(mau.Rejected, "unreadable"): + mau.parse_export("/nonexistent/export.csv") + + def test_apples_daily_active_devices_is_refused_for_the_30_day_metric(self): + with self.assertRaisesRegex(mau.Rejected, "of Active Devices: export Active Last 30 Days"): + self.parse(apple(("9/30/26", "25000.0"), column="Active Devices")) + + def test_a_blank_play_figure_is_skipped(self): + self.assertEqual(self.parse(play(("Sep 30, 2026", "", ""), ("Oct 1, 2026", "1,000", "")))[1], + {"2026-10-01": 1000}) + + def test_another_report_is_refused(self): + with self.assertRaisesRegex(mau.Rejected, "neither a Play Console MAU export"): + self.parse('Date,"Daily active users (DAU): All countries / regions"\n' + '"Sep 30, 2026","1"\n') + + def test_a_console_in_another_language_is_refused(self): + with self.assertRaisesRegex(mau.Rejected, "line 2"): + self.parse(play(("30 sept. 2026", "608 002", ""))) + + def test_a_cut_off_last_row_is_refused(self): + with self.assertRaisesRegex(mau.Rejected, "line 4"): + self.parse(ANDROID_SEPTEMBER.rsplit("\n", 2)[0] + "\nSep 3") + with self.assertRaisesRegex(mau.Rejected, "line 9"): + self.parse(IOS_SEPTEMBER + "9/3") + + def test_a_figure_cut_off_mid_number_is_refused(self): + with self.assertRaisesRegex(mau.Rejected, "not a CSV export"): + self.parse(ANDROID_SEPTEMBER.rsplit("\n", 2)[0] + '\n"Sep 30, 2026","105,0') + + def test_a_non_figure_is_refused(self): + with self.assertRaisesRegex(mau.Rejected, "for a figure"): + self.parse(play(("Sep 30, 2026", "n/a", ""))) + + +class DesktopTest(unittest.TestCase): + def test_latest_skips_prereleases(self): + self.assertIs(mau.latest(DESKTOP), LATEST) + + def test_counts_installers_only_and_adds_flathub_to_linux(self): + self.assertEqual(mau.platform_totals(LATEST, 15), + {"linux": 30, "macos": 34, "windows": 100}) + + def test_flathub_counts_from_the_release_day_on(self): + self.assertEqual(mau.flathub_installs_since(FLATHUB, "2026-07-10"), 15) + + def test_flathub_refuses_a_release_older_than_its_window(self): + with self.assertRaisesRegex(RuntimeError, "from 2026-07-09 only"): + mau.flathub_installs_since(FLATHUB, "2026-07-01") + + def test_an_error_status_fails(self): + session = FakeSession([FakeResponse({"message": "nope"}, status_code=404)]) + with self.assertRaisesRegex(RuntimeError, "HTTP 404"): + mau.desktop_downloads(session) + + +class AppleTest(unittest.TestCase): + def request(self, rid, access, stopped=False): + return {"id": rid, "attributes": {"accessType": access, + "stoppedDueToInactivity": stopped}} + + def instance(self, iid, processed): + return {"id": iid, "attributes": {"granularity": "DAILY", "processingDate": processed}} + + def page(self, *data): + return FakeResponse({"data": list(data), "links": {}}) + + def segment(self, *rows): + return GzipResponse("Date\tApp Name\tApp Apple Identifier\tDownloading Users\t" + "Users Opting-In\n" + "".join(f"{d}\tSession\t1\t{n}\t{o}\n" + for d, n, o in rows)) + + def test_the_latest_processed_instance_wins_and_older_ones_are_skipped(self): + asc = FakeSession([ + self.page(self.request("snap", "ONE_TIME_SNAPSHOT"), self.request("on", "ONGOING")), + self.page({"id": "r-snap"}), + self.page(self.instance("old", "2026-07-01"), self.instance("s", "2026-10-06")), + self.page({"id": "r-on"}), + self.page(self.instance("o", "2026-10-07")), + self.page({"attributes": {"url": "https://s3/s"}}), + self.page({"attributes": {"url": "https://s3/o"}}), + ]) + downloads = FakeSession([self.segment(("2026-09-30", 100, 25), ("2026-10-05", 10, 1)), + self.segment(("2026-10-05", 12, 3), ("2026-10-06", 8, 2))]) + days = mau.opt_in_days(asc, downloads, since="2026-08-01") + self.assertEqual(days, {"2026-09-30": (100, 25), "2026-10-05": (12, 3), + "2026-10-06": (8, 2)}) + self.assertEqual([url for _, url, _ in downloads.calls], ["https://s3/s", "https://s3/o"]) + + def test_a_stopped_ongoing_request_fails(self): + asc = FakeSession([self.page(self.request("on", "ONGOING", stopped=True))]) + with self.assertRaisesRegex(RuntimeError, "stopped .* for inactivity"): + mau.opt_in_days(asc, FakeSession([]), since="2026-08-01") + + def test_the_rate_sums_the_month_only(self): + days = {"2026-08-31": (1000, 1000), "2026-09-01": (300, 60), "2026-09-30": (100, 40)} + self.assertEqual(mau.opt_in(days, date(2026, 9, 30)), (400, 100)) + self.assertIsNone(mau.opt_in(days, date(2026, 7, 31))) + self.assertEqual(mau.ios_estimate(42000, (400, 100)), 168000) + + def test_an_empty_key_names_the_file_to_fill(self): + with tempfile.NamedTemporaryFile("w", suffix=".p8") as key, \ + mock.patch.object(mau, "desktop_downloads", return_value=DOWNLOADS), \ + mock.patch.object(mau, "apk_downloads", return_value=(ANDROID[1], 64)), \ + contextlib.redirect_stdout(io.StringIO()), \ + self.assertRaisesRegex(SystemExit, "is empty: put the App Store Connect key's " + ".p8 at /etc/session-ops/asc-key.p8"): + mau.read_sources(mock.Mock(asc_key=key.name), date(2026, 9, 30)) + + def test_apks_of_the_latest_stable_android_release_only(self): + session = FakeSession([FakeResponse(ANDROID)]) + self.assertEqual(mau.apk_downloads(session), (ANDROID[1], 64)) + + +class MergeTest(unittest.TestCase): + def empty(self): + return {"android": {}, "ios": {}, "revisions": {}} + + def test_a_later_export_wins_and_the_change_is_kept_from_the_first_figure(self): + history = {**self.empty(), "android": {"2026-09-29": 104000}} + with contextlib.redirect_stdout(io.StringIO()): + mau.merge(history, [("a", "android", {"2026-09-29": 104500, "2026-09-30": 1}), + ("b", "android", {"2026-09-29": 104600, "2026-09-30": 1})]) + self.assertEqual(history["android"], {"2026-09-29": 104600, "2026-09-30": 1}) + self.assertEqual(history["revisions"], {"android 2026-09-29": [104000, 104600]}) + + def test_a_figure_changed_back_is_no_revision(self): + history = {**self.empty(), "ios": {"2026-09-30": 42000}} + with contextlib.redirect_stdout(io.StringIO()): + mau.merge(history, [("a", "ios", {"2026-09-30": 42100}), + ("b", "ios", {"2026-09-30": 42000})]) + self.assertEqual(history["revisions"], {}) + + def test_apple_completing_a_provisional_day_is_no_revision(self): + history = {**self.empty(), "ios": {"2026-09-30": 41000.5}} + mau.merge(history, [("a", "ios", {"2026-09-30": 42000})]) + self.assertEqual(history["revisions"], {}) + + def test_only_the_days_shown_are_listed(self): + history = {**self.empty(), "revisions": {"ios 2026-09-30": [1, 2], + "android 2026-08-31": [3, 4], + "android 2026-09-15": [5, 6]}} + self.assertEqual(mau.shown_revisions(history, {"2026-09-30", "2026-08-31"}), + [("android", "2026-08-31", 3, 4), ("ios", "2026-09-30", 1, 2)]) + + +class MessageTest(unittest.TestCase): + def report(self, history=HISTORY, revisions=(), sources=SOURCES): + return mau.report_message(date(2026, 9, 30), history, list(revisions), sources) + + def test_each_store_gives_its_month_end_and_the_change_on_the_month_before(self): + message = self.report() + self.assertIn("**Monthly active users, September 2026**", message) + self.assertIn("Android, Play: **105,000** (+5,000, +5.0% on August)", message) + self.assertIn("iOS: **168,000** (−32,000, -16.0% on August), estimated", message) + + def test_ios_names_the_opt_in_report_and_its_counts(self): + self.assertIn("42,000 active in the 30 days to 30 September, scaled by September's " + "opt-in rate from Apple's App Opt In report: 25.0%, 100 of 400 " + "first-time downloaders", self.report()) + + def test_apks_and_the_fdroid_estimate_share_a_line_and_everything_adds_up(self): + message = self.report() + self.assertIn("Android, outside Play: **27,387** (GitHub APKs 64 · F-Droid 27,323, " + "estimated)", message) + self.assertIn("Desktop: **164** (Linux 30 · macOS 34 · Windows 100)", message) + self.assertIn("**Total: 300,551**", message) + self.assertIn("estimated as 10% of the other figures", message) + self.assertIn("downloads of 1.33.5's APKs since its release on 13 July", message) + self.assertIn("downloads of v1.18.1 since its release on 10 July", message) + + def test_no_change_without_the_month_befores_rate(self): + message = self.report(sources={**SOURCES, "opt_in": ((400, 100), None)}) + self.assertIn("iOS: **168,000**, estimated", message) + + def test_no_opt_in_rate_fails_rather_than_guess(self): + with self.assertRaisesRegex(RuntimeError, "no opt-in rate for September 2026"): + self.report(sources={**SOURCES, "opt_in": (None, None)}) + + def test_a_drop_is_signed(self): + history = {**HISTORY, "android": {"2026-08-31": 107000, "2026-09-30": 105000}} + self.assertIn("(−2,000, -1.9% on August)", self.report(history)) + + def test_revisions_are_listed(self): + message = self.report(revisions=[("ios", "2026-09-28", 41200, 41300)]) + self.assertIn("revised iOS 28 Sep 41,200 → 41,300", message) + + def test_the_reminder_names_each_missing_platform_and_the_upload_command(self): + message = mau.reminder_message(date(2026, 9, 30), ["android", "ios"]) + self.assertIn("Android and iOS active users for September 2026 are missing", message) + self.assertIn("- iOS: App Store Connect → Analytics", message) + self.assertIn("covering 30 September", message) + self.assertIn("`/mau-upload`", message) + + def test_the_reminder_says_when_the_last_export_was_provisional(self): + message = mau.reminder_message(date(2026, 9, 30), ["ios"], ["ios"]) + self.assertIn("covering 30 September → Export (the last export had it still being " + "counted)", message) + + def test_the_reminder_for_one_platform_is_singular(self): + message = mau.reminder_message(date(2026, 9, 30), ["ios"]) + self.assertIn("iOS active users for September 2026 is missing", message) + self.assertNotIn("Android", message) + + +class RunTest(unittest.TestCase): + def setUp(self): + self.state = tempfile.mkdtemp() + self.addCleanup(lambda: __import__("shutil").rmtree(self.state)) + os.makedirs(os.path.join(self.state, "inbox")) + patcher = mock.patch.dict(os.environ, {"MAU_DISCORD_WEBHOOK_URL": "https://hook"}) + patcher.start() + self.addCleanup(patcher.stop) + + def drop(self, name, text): + with open(os.path.join(self.state, "inbox", name), "w", encoding="utf-8") as handle: + handle.write(text) + + def run_on(self, day, *args, responses=(FakeResponse({}),)): + session = FakeSession(list(responses)) + now = datetime(*day, 12, tzinfo=mau.ZONE) + with mock.patch.object(mau, "datetime", wraps=datetime) as clock, \ + mock.patch.object(mau.http, "Session", return_value=session), \ + mock.patch.object(mau, "read_sources", return_value=SOURCES), \ + contextlib.redirect_stdout(io.StringIO()) as out: + clock.now.return_value = now + mau.main(["--state", self.state, *args]) + return session, out.getvalue() + + def posted(self, session): + return [kwargs["json"]["content"] for method, _, kwargs in session.calls + if method == "POST"] + + def history(self): + with open(os.path.join(self.state, mau.HISTORY), encoding="utf-8") as handle: + return json.load(handle) + + def listing(self, folder): + return os.listdir(os.path.join(self.state, folder)) + + def drop_both(self): + self.drop("All countries _ regions.csv", ANDROID_SEPTEMBER) + self.drop("session_private_messenger-active_last_30_days.csv", IOS_SEPTEMBER) + + def test_both_exports_post_once_and_are_filed_away(self): + self.drop_both() + session, _ = self.run_on((2026, 10, 9)) + (message,) = self.posted(session) + self.assertIn("Android, Play: **105,000**", message) + self.assertIn("iOS: **168,000**", message) + self.assertEqual(self.listing("inbox"), []) + self.assertEqual(len(self.listing("done")), 2) + self.assertEqual(self.history()["posted"], ["2026-09"]) + + session, out = self.run_on((2026, 10, 10)) + self.assertEqual(session.calls, []) + self.assertIn("2026-09 already posted", out) + + def test_one_platform_alone_waits_then_is_reminded_of_the_other(self): + self.drop("a.csv", ANDROID_SEPTEMBER) + session, out = self.run_on((2026, 10, 9)) + self.assertEqual(session.calls, []) + self.assertIn("Waiting for 2026-09-30 from ios", out) + + session, _ = self.run_on((2026, 10, 10)) + (message,) = self.posted(session) + self.assertIn("iOS active users for September 2026 is missing", message) + self.assertEqual(self.history()["posted"], []) + + self.drop("i.csv", IOS_SEPTEMBER) + session, _ = self.run_on((2026, 10, 11)) + self.assertIn("iOS: **168,000**", self.posted(session)[0]) + self.assertEqual(self.history()["posted"], ["2026-09"]) + + def test_an_export_without_the_month_end_is_kept_but_not_posted(self): + self.drop("early.csv", play(("Sep 29, 2026", "104,500", ""))) + self.drop("i.csv", IOS_SEPTEMBER) + session, _ = self.run_on((2026, 10, 10)) + self.assertIn("Android active users for September 2026 is missing", self.posted(session)[0]) + self.assertEqual(self.history()["android"], {"2026-09-29": 104500}) + + def test_a_rejected_file_is_set_aside_and_fails_the_run_after_the_rest(self): + self.drop_both() + self.drop("wrong.csv", "Date,Installs\n") + with self.assertRaisesRegex(SystemExit, "wrong.csv: neither"): + self.run_on((2026, 10, 9)) + self.assertEqual(self.listing("inbox"), []) + self.assertEqual(len(self.listing("rejected")), 1) + self.assertEqual(self.history()["posted"], ["2026-09"]) + + def test_a_dotfile_is_left_for_rsync_to_finish(self): + self.drop(".export.csv.Ab12Cd", ANDROID_SEPTEMBER) + session, _ = self.run_on((2026, 10, 9)) + self.assertEqual(session.calls, []) + self.assertEqual(self.listing("inbox"), [".export.csv.Ab12Cd"]) + + def test_a_refused_post_is_not_recorded_and_fails_the_run(self): + self.drop_both() + with self.assertRaisesRegex(RuntimeError, "Discord did not accept"): + self.run_on((2026, 10, 9), responses=[FakeResponse({}, status_code=400)]) + self.assertEqual(self.history()["posted"], []) + self.assertIn("2026-09-30", self.history()["ios"]) + + def test_a_dry_run_moves_and_writes_nothing(self): + self.drop_both() + session, out = self.run_on((2026, 10, 9), "--dry-run") + self.assertEqual(self.posted(session), []) + self.assertIn("iOS: **168,000**", out) + self.assertEqual(len(self.listing("inbox")), 2) + self.assertFalse(os.path.exists(os.path.join(self.state, mau.HISTORY))) + + def test_an_unreadable_history_stops_the_run_and_empties_the_inbox(self): + with open(os.path.join(self.state, mau.HISTORY), "w", encoding="utf-8") as handle: + handle.write("{") + self.drop("a.csv", ANDROID_SEPTEMBER) + with self.assertRaises(ValueError): + self.run_on((2026, 10, 9)) + self.assertEqual(self.listing("inbox"), []) + (held,) = self.listing("rejected") + self.assertTrue(held.endswith("-a.csv")) + with open(os.path.join(self.state, mau.HISTORY), encoding="utf-8") as handle: + self.assertEqual(handle.read(), "{") + + def test_an_unreadable_export_is_rejected_not_retried(self): + self.drop("a.csv", ANDROID_SEPTEMBER) + os.chmod(os.path.join(self.state, "inbox", "a.csv"), 0) + if os.access(os.path.join(self.state, "inbox", "a.csv"), os.R_OK): + self.skipTest("running as root, which reads any file") + with self.assertRaisesRegex(SystemExit, "a.csv: unreadable"): + self.run_on((2026, 10, 9)) + self.assertEqual(self.listing("inbox"), []) + + def test_a_provisional_ios_month_end_waits_for_a_final_one(self): + self.drop("a.csv", ANDROID_SEPTEMBER) + self.drop("i.csv", apple(("9/30/26", "41000.5"))) + session, _ = self.run_on((2026, 10, 10)) + (message,) = self.posted(session) + self.assertIn("iOS active users for September 2026 is missing", message) + self.assertIn("still being counted", message) + self.assertEqual(self.history()["posted"], []) + + self.drop("i2.csv", IOS_SEPTEMBER) + session, _ = self.run_on((2026, 10, 11)) + self.assertIn("iOS: **168,000**", self.posted(session)[0]) + + def test_revisions_from_an_earlier_run_are_posted_then_cleared(self): + self.drop("i.csv", IOS_SEPTEMBER) + self.run_on((2026, 10, 3)) + self.drop("i2.csv", apple(("8/31/26", "40100.0"), ("9/15/26", "1.0"), ("9/30/26", "42000.0"))) + self.drop("i3.csv", apple(("9/15/26", "2.0"))) + self.run_on((2026, 10, 4)) + self.drop("a.csv", ANDROID_SEPTEMBER) + session, _ = self.run_on((2026, 10, 9)) + message = self.posted(session)[0] + self.assertIn("revised iOS 31 Aug 40,000 → 40,100", message) + self.assertNotIn("15 Sep", message) + self.assertEqual(self.history()["revisions"], {})