From 1cad6a27161907bab45f2344bbad4750f9b92fd6 Mon Sep 17 00:00:00 2001 From: Audric Ackermann Date: Thu, 8 Oct 2026 17:49:37 +1100 Subject: [PATCH 1/7] feat: start jobs and hand mau its export from Discord A Discord app on webhooks.session.codes, answering two guild slash commands: - /run job: starts session-ops@.service for a job jobs.toml marks `discord = true` (crowdin-sync and crowdin-duplicates), refused while that job or any queued job is running or waiting to. - /mau-upload file: downloads the attachment, checks it with mau's own parser, and moves it into mau's inbox, whose path unit runs the job. The relay runs as opsbot behind nginx on 127.0.0.1:8081, with no gateway connection. It refuses a request that Discord did not sign, that comes from another server, or whose author is in neither allowlist. The polkit rule `session-ops polkit` writes from jobs.toml lets opsbot start the discord jobs and nothing else, and the mau inbox is the one path it can write. --- README.md | 1 + deploy/README.md | 30 ++- deploy/env/discord.env.example | 14 ++ deploy/install.sh | 40 ++-- deploy/nginx-webhooks.conf | 26 ++- deploy/session-ops-discord.service | 40 ++++ docs/jobs/mau.md | 3 +- docs/jobs/session-ops-discord.md | 28 +++ src/session_ops/jobs.toml | 3 + src/session_ops/ops/discord_commands.py | 61 +++++ src/session_ops/ops/discord_relay.py | 282 ++++++++++++++++++++++++ src/session_ops/ops/registry.py | 7 + src/session_ops/ops/runner.py | 17 ++ src/session_ops/ops/units.py | 18 ++ src/session_ops/platforms/mau.py | 5 +- tests/goldens/units/polkit.rules | 9 + tests/ops/test_deploy.py | 7 + tests/ops/test_discord_relay.py | 261 ++++++++++++++++++++++ tests/ops/test_registry.py | 14 ++ 19 files changed, 841 insertions(+), 25 deletions(-) create mode 100644 deploy/env/discord.env.example create mode 100644 deploy/session-ops-discord.service create mode 100644 docs/jobs/session-ops-discord.md create mode 100644 src/session_ops/ops/discord_commands.py create mode 100644 src/session_ops/ops/discord_relay.py create mode 100644 tests/goldens/units/polkit.rules create mode 100644 tests/ops/test_discord_relay.py diff --git a/README.md b/README.md index f720f8d..8e4aac0 100644 --- a/README.md +++ b/README.md @@ -10,6 +10,7 @@ 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) | diff --git a/deploy/README.md b/deploy/README.md index 8808cd2..e03e3e4 100644 --- a/deploy/README.md +++ b/deploy/README.md @@ -10,6 +10,7 @@ does the work: accounts, venv, env files, units, timers, and migrating an older | `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]` @@ -53,6 +54,31 @@ 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. Set its Interactions Endpoint URL to `https://webhooks.session.codes/discord/interactions`. + Discord checks it with a signed PING before saving. +4. 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 + ``` + +5. The commands start hidden from everyone. In Server Settings → Integrations, allow + them for the role or channel that runs jobs. + +Register again after changing which jobs have `discord = true`. + ## Secrets Each `/etc/session-ops/.env` has a commented `.env.example` beside it, @@ -70,6 +96,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 @@ -89,7 +117,7 @@ done 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 +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/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/install.sh b/deploy/install.sh index 53d1000..76ce2db 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,7 +68,7 @@ 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 mau; 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 @@ -92,10 +92,11 @@ 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. +# A watched job's inbox: root and the Discord relay drop files in, and the job's account +# moves them out. "$OPS" list --watched | while read -r job user dir; do install -d -o "$user" -g "$user" -m 711 "$STATE/$job" - install -d -o "$user" -g "$user" -m 700 "$dir" + install -d -o "$user" -g opsbot -m 770 "$dir" done install -m 644 "$ROOT/deploy/session-ops.tmpfiles" /etc/tmpfiles.d/session-ops.conf @@ -115,6 +116,12 @@ install -m 644 "$ROOT"/deploy/*.service "$ROOT"/deploy/*.timer "$ROOT"/deploy/*. rm -f "$UNITS"/session-ops@*.service.d/job.conf "$UNITS"/session-ops@*.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) @@ -173,16 +180,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/ diff --git a/docs/jobs/session-ops-discord.md b/docs/jobs/session-ops-discord.md new file mode 100644 index 0000000..c211cf1 --- /dev/null +++ b/docs/jobs/session-ops-discord.md @@ -0,0 +1,28 @@ +# 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. The job posts its outcome in its own channel. | +| `/mau-upload file:` | Downloads the attachment, checks it parses as the Play Console 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, and the server's Integrations settings show them to the right +role; that only hides them, and these checks are what refuse everyone else. + +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/jobs.toml b/src/session_ops/jobs.toml index 179c1ba..964f203 100644 --- a/src/session_ops/jobs.toml +++ b/src/session_ops/jobs.toml @@ -10,6 +10,7 @@ # 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 # @@ -73,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" @@ -113,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"] diff --git a/src/session_ops/ops/discord_commands.py b/src/session_ops/ops/discord_commands.py new file mode 100644 index 0000000..e0dff09 --- /dev/null +++ b/src/session_ops/ops/discord_commands.py @@ -0,0 +1,61 @@ +"""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; +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. +""" +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 MAU export", + "options": [{"type": OPTION_ATTACHMENT, "name": "file", "required": True, + "description": "The CSV the saved MAU report exports"}]}, + ] + + +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..e9c96c7 --- /dev/null +++ b/src/session_ops/ops/discord_relay.py @@ -0,0 +1,282 @@ +"""Discord slash commands for session-ops jobs: POST /discord/interactions. + + /run job: starts session-ops@.service + /mau-upload file: puts a Play Console 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 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"}) + +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) + + +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 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: + blocking = busy(job) + started = None if blocking else systemctl("start", "--no-block", unit) + except (OSError, subprocess.SubprocessError, RuntimeError) 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 blocking: + return reply(f"⏳ `{blocking}` is running or waiting to; try again once it has finished.") + 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"Play Console exports 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: + 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 export, {len(days)} days from {min(days)} to {max(days)}, " + f"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 7085af8..6d06a06 100644 --- a/src/session_ops/ops/registry.py +++ b/src/session_ops/ops/registry.py @@ -11,6 +11,8 @@ "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 @dataclass(frozen=True) @@ -31,6 +33,7 @@ class Job: queued: bool = False after: tuple = () watch: str = None + discord: bool = False @property def scheduled(self): @@ -93,6 +96,10 @@ def load(path=REGISTRY): 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") + 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 7184c02..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 @@ -243,6 +245,11 @@ def main(argv=None): 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) @@ -270,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 1bdbf45..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): @@ -47,6 +50,21 @@ def dropins(jobs, queue): 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 index ba501c4..d9615df 100644 --- a/src/session_ops/platforms/mau.py +++ b/src/session_ops/platforms/mau.py @@ -5,6 +5,7 @@ session-ops run mau [--dry-run] rsync "All countries _ regions.csv" root@:/var/lib/session-ops/mau/inbox/ + /mau-upload file:.csv # in Discord, through session-ops-discord 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; @@ -209,8 +210,8 @@ def reminder_message(month_end, inbox): return "\n".join([ f"⏰ **Android MAU for {month_end:%B %Y} is missing.**", "In Play Console, open Statistics → Saved reports → the MAU report, check it covers " - f"{month_end:%-d %B}, and Export report → CSV. Then copy it to the inbox, and the " - "figures post as soon as it lands:", + f"{month_end:%-d %B}, and Export report → CSV. Then hand it over with `/mau-upload` " + "here, or copy it to the inbox, and the figures post as soon as it lands:", f'`rsync ".csv" {os.environ.get("MAU_INBOX_HOST") or socket.getfqdn()}:{inbox}/`', ]) 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 c66b8f6..850140d 100644 --- a/tests/ops/test_deploy.py +++ b/tests/ops/test_deploy.py @@ -89,6 +89,13 @@ 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([])) + 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..a4c006f --- /dev/null +++ b/tests/ops/test_discord_relay.py @@ -0,0 +1,261 @@ +""" + 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.MAU_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): + self.running, self.waiting, self.start_error = set(running), set(waiting), start_error + self.started = [] + + 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] == "start": + self.started.append(args[-1]) + 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_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), {"2026-09-29": 1200, "2026-09-30": 1250}) + self.assertEqual(os.stat(path).st_mode & 0o777, 0o644) + self.assertIn("2 days from 2026-09-29 to 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"), "exports 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 f7e82b5..1c61072 100644 --- a/tests/ops/test_registry.py +++ b/tests/ops/test_registry.py @@ -80,6 +80,20 @@ 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_the_shipped_registry_loads(self): names = [job.name for job in registry.load()] self.assertIn("zendesk-digest", names) From 745285977c63ba54b1d57f9eb3888b47c524f96d Mon Sep 17 00:00:00 2001 From: Audric Ackermann Date: Thu, 8 Oct 2026 20:51:44 +1100 Subject: [PATCH 2/7] fix: keep a /run job clear of the queue, and two /run calls apart /run refused while the queue was busy, but nothing stopped the queue's timer starting while a /run job was still going: /run crowdin-duplicates at 09:55 then ran beside the 10:00 crowdin-sync, two Crowdin clients against one rate limit. /run now also refuses when the queue timer's next run is sooner than now plus the job's timeout, since systemd kills the job by then. The next run comes from `systemctl list-timers --output=json`: on systemd 252, the host's version, `show --timestamp=unix -P NextElapseUSecRealtime` still prints local time. The check and `systemctl start --no-block` now run under one lock, so a second /run a moment later sees the first's job queued. A job's timeout is parsed as a systemd time span when jobs.toml loads, so a bad one fails there. --- src/session_ops/ops/discord_relay.py | 44 ++++++++++++++++++++++++---- src/session_ops/ops/registry.py | 31 ++++++++++++++++++++ tests/ops/test_discord_relay.py | 33 +++++++++++++++++++-- tests/ops/test_registry.py | 10 +++++++ 4 files changed, 111 insertions(+), 7 deletions(-) diff --git a/src/session_ops/ops/discord_relay.py b/src/session_ops/ops/discord_relay.py index e9c96c7..b555561 100644 --- a/src/session_ops/ops/discord_relay.py +++ b/src/session_ops/ops/discord_relay.py @@ -21,6 +21,7 @@ import json import os import subprocess +import threading import time from urllib.parse import urlsplit @@ -44,6 +45,7 @@ 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 @@ -54,6 +56,8 @@ NO_PINGS = {"parse": []} app = FastAPI(docs_url=None, redoc_url=None, openapi_url=None) +# Held across the check and the start: two /run a moment apart would both pass the check. +_starting = threading.Lock() def env(name, default=None): @@ -147,6 +151,35 @@ def busy(job): 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) @@ -154,14 +187,15 @@ def handle_run(interaction): return reply(f"❌ `{name}` is not a job /run can start.") unit = unit_name(job.name) try: - blocking = busy(job) - started = None if blocking else systemctl("start", "--no-block", unit) - except (OSError, subprocess.SubprocessError, RuntimeError) as exc: + 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 blocking: - return reply(f"⏳ `{blocking}` is running or waiting to; try again once it has finished.") + 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) diff --git a/src/session_ops/ops/registry.py b/src/session_ops/ops/registry.py index 6d06a06..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 @@ -13,6 +14,28 @@ 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) @@ -43,6 +66,10 @@ 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) @@ -98,6 +125,10 @@ def load(path=REGISTRY): 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/tests/ops/test_discord_relay.py b/tests/ops/test_discord_relay.py index a4c006f..9b70d61 100644 --- a/tests/ops/test_discord_relay.py +++ b/tests/ops/test_discord_relay.py @@ -32,9 +32,10 @@ class FakeSystemd: - def __init__(self, running=(), waiting=(), start_error=None): + 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.started = [] + self.queue_at = queue_at + self.started, self.locked = [], [] def __call__(self, *args): out, code, err = "", 0, "" @@ -43,8 +44,12 @@ def __call__(self, *args): 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) @@ -156,6 +161,30 @@ def test_a_running_or_waiting_queue_job_holds_it_off(self): 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): diff --git a/tests/ops/test_registry.py b/tests/ops/test_registry.py index 1c61072..cd1a94e 100644 --- a/tests/ops/test_registry.py +++ b/tests/ops/test_registry.py @@ -94,6 +94,16 @@ def test_more_jobs_than_discord_offers_choices_for_is_refused(self): 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) From 8f0858bbed06d8a548c8cbed12597f987b195802 Mon Sep 17 00:00:00 2001 From: Audric Ackermann Date: Thu, 8 Oct 2026 20:52:02 +1100 Subject: [PATCH 3/7] docs: copy the Discord route into the live nginx file before setting the URL Discord checks the Interactions Endpoint URL with a request that, without the route, hits the catch-all 404, and it then refuses to save the URL. --- deploy/README.md | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/deploy/README.md b/deploy/README.md index e03e3e4..661fa29 100644 --- a/deploy/README.md +++ b/deploy/README.md @@ -64,9 +64,14 @@ What `/run` and `/mau-upload` do, and who may run them: 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. Set its Interactions Endpoint URL to `https://webhooks.session.codes/discord/interactions`. +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. -4. Register the commands, with the bot token from the Bot page, typed in rather than kept: +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 @@ -74,7 +79,7 @@ What `/run` and `/mau-upload` do, and who may run them: /opt/session-ops/.venv/bin/session-ops discord-register ``` -5. The commands start hidden from everyone. In Server Settings → Integrations, allow +6. The commands start hidden from everyone. In Server Settings → Integrations, allow them for the role or channel that runs jobs. Register again after changing which jobs have `discord = true`. From 0f1665ca186861e2f5572a2a4269419e54f37bcb Mon Sep 17 00:00:00 2001 From: Audric Ackermann Date: Thu, 8 Oct 2026 20:52:13 +1100 Subject: [PATCH 4/7] docs: the hidden commands still show to server administrators default_member_permissions "0" hides a command from every member without Administrator; the owner and administrators still see it. The relay's allowlist refuses them like anyone else. --- deploy/README.md | 5 +++-- docs/jobs/session-ops-discord.md | 5 +++-- src/session_ops/ops/discord_commands.py | 7 ++++--- 3 files changed, 10 insertions(+), 7 deletions(-) diff --git a/deploy/README.md b/deploy/README.md index 661fa29..6731ea0 100644 --- a/deploy/README.md +++ b/deploy/README.md @@ -79,8 +79,9 @@ What `/run` and `/mau-upload` do, and who may run them: /opt/session-ops/.venv/bin/session-ops discord-register ``` -6. The commands start hidden from everyone. In Server Settings → Integrations, allow - them for the role or channel that runs jobs. +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`. diff --git a/docs/jobs/session-ops-discord.md b/docs/jobs/session-ops-discord.md index c211cf1..fb16e4b 100644 --- a/docs/jobs/session-ops-discord.md +++ b/docs/jobs/session-ops-discord.md @@ -20,8 +20,9 @@ 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, and the server's Integrations settings show them to the right -role; that only hides them, and these checks are what refuse everyone else. +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 diff --git a/src/session_ops/ops/discord_commands.py b/src/session_ops/ops/discord_commands.py index e0dff09..2f2f0b5 100644 --- a/src/session_ops/ops/discord_commands.py +++ b/src/session_ops/ops/discord_commands.py @@ -3,9 +3,10 @@ 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; -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. +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 0e37ae77e231c1920b57d97349fd8ff7443c25f7 Mon Sep 17 00:00:00 2001 From: Audric Ackermann Date: Thu, 8 Oct 2026 20:52:32 +1100 Subject: [PATCH 5/7] fix: give the Discord relay's group to mau's inbox alone install.sh made every watched job's inbox writable by opsbot, though the relay writes only mau's. Other watched inboxes stay the job account's alone, at 0700. A test ties the relay unit's ReadWritePaths= and install.sh to mau's watch. --- deploy/install.sh | 10 +++++++--- tests/ops/test_deploy.py | 12 +++++++++++- 2 files changed, 18 insertions(+), 4 deletions(-) diff --git a/deploy/install.sh b/deploy/install.sh index 76ce2db..8d89629 100755 --- a/deploy/install.sh +++ b/deploy/install.sh @@ -92,11 +92,15 @@ 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 and the Discord relay drop files in, and the job's account -# moves them out. +# 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" - install -d -o "$user" -g opsbot -m 770 "$dir" + 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 diff --git a/tests/ops/test_deploy.py b/tests/ops/test_deploy.py index 850140d..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__)))) @@ -96,6 +96,16 @@ 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") From 8f0e984f1a034e0767c0706b179b373a70a14d63 Mon Sep 17 00:00:00 2001 From: Audric Ackermann Date: Thu, 8 Oct 2026 20:54:19 +1100 Subject: [PATCH 6/7] docs: /run also waits out the queue's next run The doc page named only the busy check. The start lock holds only within one process, which is how session-ops-discord.service runs uvicorn. --- docs/jobs/session-ops-discord.md | 2 +- src/session_ops/ops/discord_relay.py | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/docs/jobs/session-ops-discord.md b/docs/jobs/session-ops-discord.md index fb16e4b..168e8eb 100644 --- a/docs/jobs/session-ops-discord.md +++ b/docs/jobs/session-ops-discord.md @@ -9,7 +9,7 @@ | 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. The job posts its outcome in its own channel. | +| `/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 the Play Console 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 diff --git a/src/session_ops/ops/discord_relay.py b/src/session_ops/ops/discord_relay.py index b555561..1df6e3d 100644 --- a/src/session_ops/ops/discord_relay.py +++ b/src/session_ops/ops/discord_relay.py @@ -56,7 +56,7 @@ NO_PINGS = {"parse": []} app = FastAPI(docs_url=None, redoc_url=None, openapi_url=None) -# Held across the check and the start: two /run a moment apart would both pass the check. +# Across the check and the start, so two /run a moment apart cannot both pass; one process only. _starting = threading.Lock() From 7292f7acc4bf9586de737a77dd587b9683cde264 Mon Sep 17 00:00:00 2001 From: Audric Ackermann Date: Thu, 8 Oct 2026 22:14:16 +1100 Subject: [PATCH 7/7] fix: start the queue at its scheduled time, with no random delay /run refuses a job whose timeout would carry it past the queue's next run, read from list-timers. With RandomizedDelaySec, re-arming the timer draws a new delay, so the queue could start up to two minutes before the time /run read. On one host the delay spread nothing. --- deploy/session-ops-queue.timer | 2 +- docs/jobs/github-prs-digest.md | 5 ++--- src/session_ops/github_prs/digest.py | 2 +- 3 files changed, 4 insertions(+), 5 deletions(-) diff --git a/deploy/session-ops-queue.timer b/deploy/session-ops-queue.timer index fd5ff9f..8448ff1 100644 --- a/deploy/session-ops-queue.timer +++ b/deploy/session-ops-queue.timer @@ -5,7 +5,7 @@ Documentation=https://github.com/session-foundation/session-shared-scripts [Timer] Persistent=yes -RandomizedDelaySec=2min +# No random delay: /run in Discord refuses a job its timeout would carry past the next run. [Install] WantedBy=timers.target diff --git a/docs/jobs/github-prs-digest.md b/docs/jobs/github-prs-digest.md index c40dae6..846af26 100644 --- a/docs/jobs/github-prs-digest.md +++ b/docs/jobs/github-prs-digest.md @@ -49,9 +49,8 @@ cache rather than something to back up. ### Late, not lost -Two runs can be further apart than 72 hours: April's DST weekend is 73, the timer's -`RandomizedDelaySec` adds up to two minutes, and a host that was down runs once when it -comes back. So the state also keeps `covered_until`, the time the last run whose every +Two runs can be further apart than 72 hours: April's DST weekend is 73, and a host that +was down runs once when it comes back. So the state also keeps `covered_until`, the time the last run whose every message Discord accepted *started* its search, and each run reaches back to whichever is earlier, that or 72 hours ago. A run whose post fails partway, or that fails before posting, never moves it forward; a dry run writes nothing. The header then names the 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. """