mirror of
https://github.com/DarkflameUniverse/DarkflameServer.git
synced 2026-10-02 10:53:44 +00:00
feat: web dashboard and playground work
The NexusDashboard-parity dashboard (dDashboardServer) and everything built on it on the experimental branch: accounts, characters, properties and moderation tools, permissions shared with in-game slash commands, economy reports, World 3D and property 3D views with client scenery, scheduled events (features, vanity changes, live events, announcements, restarts), vanity files and events, the CDClient browser, the message inspector with saved captures, chat filter tools, community challenges, live ops, the AI moderator helper, and the server-side changes they need. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
4
tools/chat-bridge/.gitignore
vendored
Normal file
4
tools/chat-bridge/.gitignore
vendored
Normal file
@@ -0,0 +1,4 @@
|
||||
__pycache__/
|
||||
config.json
|
||||
*.state.json
|
||||
*.state.json.tmp
|
||||
172
tools/chat-bridge/README.md
Normal file
172
tools/chat-bridge/README.md
Normal file
@@ -0,0 +1,172 @@
|
||||
# Chat bridge (example)
|
||||
|
||||
An example program that connects the game's chat to another chat service through the dashboard's public chat API.
|
||||
Discord is the first adapter; a console adapter (read/type chat in a terminal) and a send-only JSON webhook adapter
|
||||
show how to add others (Matrix, IRC, Slack, ...). The server knows nothing about any of them: the bridge uses only the
|
||||
documented API (`docs/Dashboard.md`, "Chat log and chat bridges" and "API"), so you can replace it with your own.
|
||||
|
||||
What it does:
|
||||
|
||||
- **Game to the other side.** Reads chat with `GET /api/chat?after=<last id>`. It keeps a WebSocket open to `/ws`
|
||||
(subscribed to `chat_message`) only to know *when* to read, so a message is never lost or sent twice: when the
|
||||
socket drops it polls every few seconds, and after a restart it carries on from the last id it saved
|
||||
(`state_file`). On its first run it starts from the newest message rather than replaying the log.
|
||||
- **The other side to the game.** Posts with `POST /api/chat/send` as `{"message", "name", "label"}`, optionally
|
||||
only in one zone (`send_zone`) or one world of it (`send_instance`). Players see `[Discord] Alice: hi`.
|
||||
- **Only zone chat by default.** Whispers and team chat are never relayed unless a link lists them *and* sets
|
||||
`"allow_private": true`. When a link wants one channel, the bridge asks the server for only that channel, so other
|
||||
rows (like whispers, which a GM 8 account can read) don't even leave the server.
|
||||
- **Nothing the chat filter stopped** (`blocked`) is relayed: no player saw it, so no one else should.
|
||||
- **No loops.** It skips `web` messages with its own label or sent by its own account, and the Discord adapter
|
||||
skips bots and webhooks (including its own).
|
||||
- **Rate limits both ways**, per link (a token bucket; `0` turns it off). Game messages over the limit are dropped
|
||||
and the next one says how many were skipped; Discord messages over the limit aren't posted in the game.
|
||||
- Messages into the game are made to fit the API: one line, up to 300 characters, names up to 32 without brackets.
|
||||
Discord mentions, channels and custom emoji are spelled out (`@bob`, `#channel`, `:lego:`), and attachments show
|
||||
as `[attachment]`. Nothing said in the game can ping anyone on Discord (`allowed_mentions` is empty) and its
|
||||
Markdown is escaped.
|
||||
|
||||
## Requirements
|
||||
|
||||
Python 3.10+ and, for live updates and Discord, [`websockets`](https://pypi.org/project/websockets/)
|
||||
(`pip install -r requirements.txt`, or your distribution's `python3-websockets`). Everything else is the standard
|
||||
library. Without `websockets` the bridge still works with the console and webhook adapters, by polling.
|
||||
|
||||
Python was chosen because the repository's tooling (`tests/smoke`) already uses it, it runs anywhere the server does,
|
||||
and with one small, widely packaged dependency there's nothing to build.
|
||||
|
||||
## Dashboard account and API token
|
||||
|
||||
1. Make an account for the bot (**Create Account** on the Accounts page), not a staff member's own: messages from
|
||||
the bot's account are ignored as its own, and its token can be revoked on its own.
|
||||
2. Give it a GM level with these permissions (see the **Permissions** page):
|
||||
- `chat_view` (read chat; GM 3+ by default),
|
||||
- `chat_send` (post in the game; GM 8+ by default),
|
||||
- `api_access` (use API tokens; everyone by default).
|
||||
|
||||
With the default levels that means GM 8. GM 8 also has `chat_private`, so the account *can* read whispers and
|
||||
team chat; this bridge still won't relay them unless you turn them on (see above).
|
||||
3. Sign in to the dashboard as the bot, open its account page and **Generate** a token on the **API Token** card
|
||||
(1 to 365 days). Copy it: it is shown once. Or through the API:
|
||||
|
||||
```sh
|
||||
curl -s -X POST http://localhost:2006/api/auth/login -H 'X-Requested-With: x' \
|
||||
-H 'Content-Type: application/json' -d '{"username":"chatbot","password":"..."}' # -> {"token": ...}
|
||||
curl -s -X POST http://localhost:2006/api/auth/token -H 'Authorization: Bearer <login token>' \
|
||||
-H 'Content-Type: application/json' -d '{"days":365}' # -> {"token": ...}
|
||||
```
|
||||
|
||||
To revoke the token, sign the account out everywhere (its account page, or
|
||||
`POST /api/accounts/<id>/sessions/revoke`); that ends every session and token of that account, which is another
|
||||
reason to give the bot its own.
|
||||
|
||||
Check it: `DLU_API_TOKEN=... python3 -m chat_bridge --config config.json --check`
|
||||
|
||||
## Discord
|
||||
|
||||
1. In the [Discord developer portal](https://discord.com/developers/applications), make an application with a bot,
|
||||
and under **Bot** turn on **Message Content Intent** (without it the bot can't read what people write).
|
||||
2. Invite it to your server with the `bot` scope and the **View Channel**, **Read Message History** and
|
||||
**Send Messages** permissions in the channel you want to bridge.
|
||||
3. Copy the bot token and the channel ID (Discord's developer mode, right-click the channel, **Copy Channel ID**).
|
||||
4. Optional: make a webhook in that channel (channel settings, **Integrations**, **Webhooks**) and set it as
|
||||
`webhook_url`. Game messages then show under each player's name (`Bob (Nimbus Station)`) instead of as the bot.
|
||||
With only a `webhook_url` and no bot token the link is one way (game to Discord).
|
||||
|
||||
## Config
|
||||
|
||||
Copy `config.example.json` to `config.json`. Any value written as `"$NAME"` is read from the environment, and
|
||||
`DLU_DASHBOARD_URL` / `DLU_API_TOKEN` override the dashboard settings, so secrets needn't be in the file.
|
||||
|
||||
```json
|
||||
{
|
||||
"dashboard": { "url": "http://localhost:2006", "token": "$DLU_API_TOKEN" },
|
||||
"state_file": "chat-bridge.state.json",
|
||||
"poll_seconds": 5,
|
||||
"resync_seconds": 60,
|
||||
"links": [
|
||||
{
|
||||
"adapter": "discord",
|
||||
"label": "Discord",
|
||||
"discord": { "bot_token": "$DISCORD_BOT_TOKEN", "channel_id": "123456789012345678", "webhook_url": "$DISCORD_WEBHOOK_URL" },
|
||||
"game": { "channels": ["zone"], "zones": [], "instances": [], "send_zone": 0 },
|
||||
"rate": { "to_game_per_minute": 20, "to_remote_per_minute": 60 }
|
||||
}
|
||||
]
|
||||
}
|
||||
```
|
||||
|
||||
- `state_file`: where the last message id is kept (relative to the config file).
|
||||
- `poll_seconds`: how often to read while the WebSocket is down; `resync_seconds`: how often to read anyway while it's up.
|
||||
- `links`: one per bridged room. Several links share one connection to the dashboard, e.g. a Discord channel for all
|
||||
zone chat and a webhook that gets only Nimbus Station's.
|
||||
- `adapter`: `discord`, `console` or `webhook`; its settings go under a key of the same name.
|
||||
- `label`: what players see in brackets (up to 16 letters, digits, spaces, `-`, `_`; default the adapter's name).
|
||||
- `game.channels`: `zone` (default), `web` (messages posted from the dashboard or other bridges), and only with
|
||||
`"allow_private": true`, `whisper` and `team`.
|
||||
- `game.zones` / `game.instances`: only relay chat from these zone IDs / world instance IDs (empty: all).
|
||||
- `game.send_zone` / `game.send_instance`: post incoming messages only in this zone / world (0 / absent: everywhere).
|
||||
- `rate`: messages per minute each way (bursts of 5).
|
||||
|
||||
Adapter settings:
|
||||
|
||||
| Adapter | Settings |
|
||||
|-----------|----------|
|
||||
| `discord` | `bot_token`, `channel_id` (to read and send); `webhook_url` (optional, to send as each player) |
|
||||
| `console` | `name` (who lines without `Name: ` are from), `read_stdin` (default true) |
|
||||
| `webhook` | `url`, `headers` (optional). Send only: POSTs `{"text", "sender", "message", "channel", "zone_id", "zone_name", "instance_id", "id"}` |
|
||||
|
||||
## Running
|
||||
|
||||
```sh
|
||||
cd tools/chat-bridge
|
||||
pip install -r requirements.txt
|
||||
DLU_API_TOKEN=... DISCORD_BOT_TOKEN=... python3 -m chat_bridge --config config.json
|
||||
```
|
||||
|
||||
Try it in a terminal first with `config.console.json`: game chat is printed, and each line you type
|
||||
(`Alice: hello` or just `hello`) is posted in the game as `[Console] Alice: hello`.
|
||||
|
||||
As a service: `chat-bridge.service` is a systemd unit; copy it to `/etc/systemd/system/`, fix the paths and user,
|
||||
put the secrets in `/etc/chat-bridge.env` (mode 600), then `systemctl enable --now chat-bridge`.
|
||||
`journalctl -u chat-bridge` shows its log.
|
||||
|
||||
## Adding a service
|
||||
|
||||
Write `chat_bridge/adapters/<name>.py` with a subclass of `Adapter` (`adapters/base.py`):
|
||||
|
||||
- `send(msg, skipped)`: post one game message (a `GameMessage`: `sender_name`, `message`, `zone_name`, ...).
|
||||
- `run(deliver)`: receive messages and `await deliver(RemoteMessage(author, text))` for each; ignore your own posts.
|
||||
A send-only adapter doesn't need it.
|
||||
|
||||
Then add it to `ADAPTERS` in `adapters/__init__.py`. `console.py` is the smallest two-way example and `webhook.py` the
|
||||
smallest one-way one. Filtering, loop prevention, rate limits and resuming are done by the bridge.
|
||||
|
||||
## Privacy
|
||||
|
||||
- Players can't see that their zone chat is being relayed. Tell them (rules, a message of the day, the news).
|
||||
- Whispers and team chat are private; relaying them needs two explicit settings. Don't.
|
||||
- The bridge only sends out what players could already see in that zone. Messages the chat filter stopped aren't
|
||||
relayed.
|
||||
- The game logs what the bridge posts (the Chat Log's `web` channel), with the bridge's account.
|
||||
- Keep the API token and bot token secret: the API token acts as the bot's account with everything its GM level can
|
||||
do, not only chat. Use a dedicated account with the lowest level that has the permissions above.
|
||||
|
||||
## Limits
|
||||
|
||||
- Messages posted in the game go through the same chat filter as players' (`chat_bridge_filter`, on by default). If
|
||||
the filter stops one, no player sees it; the API only answers with a request id, so the bridge can't tell the
|
||||
Discord user. (The Chat Log does not mark such web messages as stopped.)
|
||||
- Live messages arrive with the dashboard's update tick (`broadcast_interval`, 2 seconds by default).
|
||||
- One Discord channel per link; threads, edits, deletions and reactions aren't relayed.
|
||||
- The rate limits drop messages rather than queueing them for long; at most 100 wait per link.
|
||||
- A message can be relayed twice only if the bridge is killed between sending a message and saving the state file.
|
||||
|
||||
## Tests
|
||||
|
||||
```sh
|
||||
cd tools/chat-bridge && python3 -m unittest discover -s tests
|
||||
```
|
||||
|
||||
Standard library only (no network): relay rules, loop prevention, formatting (game and Discord), resuming from the
|
||||
last id, paging, rate limits and config parsing.
|
||||
25
tools/chat-bridge/chat-bridge.service
Normal file
25
tools/chat-bridge/chat-bridge.service
Normal file
@@ -0,0 +1,25 @@
|
||||
# systemd unit for the chat bridge. Copy to /etc/systemd/system/, adjust the paths and user, then:
|
||||
# sudo systemctl daemon-reload && sudo systemctl enable --now chat-bridge
|
||||
[Unit]
|
||||
Description=DarkflameServer chat bridge
|
||||
After=network-online.target
|
||||
Wants=network-online.target
|
||||
|
||||
[Service]
|
||||
Type=simple
|
||||
User=dlu
|
||||
WorkingDirectory=/opt/DarkflameServer/tools/chat-bridge
|
||||
# DLU_API_TOKEN=..., DISCORD_BOT_TOKEN=..., DISCORD_WEBHOOK_URL=... (chmod 600)
|
||||
EnvironmentFile=/etc/chat-bridge.env
|
||||
ExecStart=/usr/bin/python3 -m chat_bridge --config /opt/DarkflameServer/tools/chat-bridge/config.json
|
||||
# No stdin: the console adapter isn't for services
|
||||
StandardInput=null
|
||||
Restart=on-failure
|
||||
RestartSec=10
|
||||
NoNewPrivileges=true
|
||||
ProtectSystem=strict
|
||||
ReadWritePaths=/opt/DarkflameServer/tools/chat-bridge
|
||||
PrivateTmp=true
|
||||
|
||||
[Install]
|
||||
WantedBy=multi-user.target
|
||||
1
tools/chat-bridge/chat_bridge/__init__.py
Normal file
1
tools/chat-bridge/chat_bridge/__init__.py
Normal file
@@ -0,0 +1 @@
|
||||
"""An example chat bridge for the DarkflameServer dashboard's chat API (see ../README.md)."""
|
||||
84
tools/chat-bridge/chat_bridge/__main__.py
Normal file
84
tools/chat-bridge/chat_bridge/__main__.py
Normal file
@@ -0,0 +1,84 @@
|
||||
"""python3 -m chat_bridge --config config.json [--check] [-v]"""
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import asyncio
|
||||
import logging
|
||||
import os
|
||||
import sys
|
||||
|
||||
from .adapters import ADAPTERS
|
||||
from .bridge import Bridge, Link
|
||||
from .config import load_config
|
||||
from .dashboard import ApiError, Dashboard
|
||||
|
||||
|
||||
class _Subscribed(Exception):
|
||||
pass
|
||||
|
||||
|
||||
def _stop() -> None:
|
||||
raise _Subscribed
|
||||
|
||||
|
||||
async def check(dashboard: Dashboard) -> int:
|
||||
"""Try the token: who it signs in as, and whether it can read and subscribe (it doesn't post anything)."""
|
||||
me = await dashboard.me()
|
||||
if not me.get("valid"):
|
||||
print("token: not valid")
|
||||
return 1
|
||||
print(f"token: {me['username']} (account {me['accountId']}, GM {me['gmLevel']})")
|
||||
ok = True
|
||||
try:
|
||||
page = await dashboard.chat_page({"limit": "1"})
|
||||
print(f"read chat (chat_view): yes, newest id {page.get('last_id')}")
|
||||
except ApiError as e:
|
||||
print(f"read chat (chat_view): no ({e})")
|
||||
ok = False
|
||||
try:
|
||||
await asyncio.wait_for(dashboard.watch_chat(lambda _e: None, _stop), 10)
|
||||
except _Subscribed:
|
||||
print("live chat (WebSocket): yes")
|
||||
except ImportError:
|
||||
print("live chat (WebSocket): no, the websockets package isn't installed (the bridge will poll)")
|
||||
except Exception as e:
|
||||
print(f"live chat (WebSocket): no ({e}); the bridge will poll")
|
||||
print("send chat (chat_send): not tried here; see the Permissions page")
|
||||
return 0 if ok else 1
|
||||
|
||||
|
||||
def main() -> int:
|
||||
parser = argparse.ArgumentParser(prog="chat_bridge", description="Bridge game chat to Discord and other services through the dashboard API")
|
||||
parser.add_argument("--config", default=os.environ.get("CHAT_BRIDGE_CONFIG", "config.json"))
|
||||
parser.add_argument("--check", action="store_true", help="check the token and permissions, then exit")
|
||||
parser.add_argument("-v", "--verbose", action="store_true")
|
||||
args = parser.parse_args()
|
||||
logging.basicConfig(level=logging.DEBUG if args.verbose else logging.INFO, stream=sys.stderr,
|
||||
format="%(asctime)s %(levelname)s %(name)s: %(message)s")
|
||||
# websockets' debug log prints request headers, which carry the API token
|
||||
logging.getLogger("websockets").setLevel(logging.INFO)
|
||||
|
||||
config = load_config(args.config, {name: cls.default_label for name, cls in ADAPTERS.items()})
|
||||
if not os.path.isabs(config.state_file):
|
||||
config.state_file = os.path.join(os.path.dirname(os.path.abspath(args.config)), config.state_file)
|
||||
dashboard = Dashboard(config.url, config.token)
|
||||
if args.check:
|
||||
return asyncio.run(check(dashboard))
|
||||
|
||||
links = []
|
||||
for link_cfg in config.links:
|
||||
if link_cfg.adapter not in ADAPTERS:
|
||||
parser.error(f"unknown adapter {link_cfg.adapter!r}; known: {', '.join(ADAPTERS)}")
|
||||
links.append(Link(link_cfg, ADAPTERS[link_cfg.adapter](link_cfg.settings)))
|
||||
try:
|
||||
asyncio.run(Bridge(config, dashboard, links).run())
|
||||
except KeyboardInterrupt:
|
||||
pass
|
||||
except RuntimeError as e:
|
||||
logging.error("%s", e)
|
||||
return 1
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
12
tools/chat-bridge/chat_bridge/adapters/__init__.py
Normal file
12
tools/chat-bridge/chat_bridge/adapters/__init__.py
Normal file
@@ -0,0 +1,12 @@
|
||||
"""The chat services the bridge can talk to. Add a new one here (and a module next to these)."""
|
||||
from .base import Adapter
|
||||
from .console import ConsoleAdapter
|
||||
from .discord import DiscordAdapter
|
||||
from .webhook import WebhookAdapter
|
||||
|
||||
ADAPTERS: dict[str, type[Adapter]] = {
|
||||
ConsoleAdapter.name: ConsoleAdapter,
|
||||
DiscordAdapter.name: DiscordAdapter,
|
||||
WebhookAdapter.name: WebhookAdapter,
|
||||
# "matrix": MatrixAdapter, "irc": IrcAdapter, ...
|
||||
}
|
||||
37
tools/chat-bridge/chat_bridge/adapters/base.py
Normal file
37
tools/chat-bridge/chat_bridge/adapters/base.py
Normal file
@@ -0,0 +1,37 @@
|
||||
"""The adapter interface: everything specific to one chat service lives behind it.
|
||||
|
||||
To add a service (Matrix, IRC, ...): subclass Adapter, implement run() and send(), and register it in
|
||||
adapters/__init__.py. The bridge handles the dashboard, filtering, loop prevention and rate limits.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Awaitable, Callable
|
||||
|
||||
from ..core import GameMessage, RemoteMessage
|
||||
|
||||
# Given to run(): call it with each message the service receives, to post it in the game
|
||||
Deliver = Callable[[RemoteMessage], Awaitable[None]]
|
||||
|
||||
|
||||
class Adapter:
|
||||
#: The config key under a link that holds this adapter's settings, e.g. "discord"
|
||||
name = "base"
|
||||
#: The default label players see ("[Discord] Bob: hi"); a link's "label" overrides it
|
||||
default_label = "Bridge"
|
||||
|
||||
def __init__(self, config: dict):
|
||||
self.config = config
|
||||
|
||||
async def run(self, deliver: Deliver) -> None:
|
||||
"""Receive messages from the service for as long as the bridge runs, calling deliver() for each.
|
||||
|
||||
Ignore messages the adapter itself posted (e.g. Discord bot/webhook messages) so they don't loop back.
|
||||
An adapter that only sends can just return.
|
||||
"""
|
||||
|
||||
async def send(self, msg: GameMessage, skipped: int = 0) -> None:
|
||||
"""Post one game chat message on the service. `skipped` is how many were dropped by the rate limit before it."""
|
||||
raise NotImplementedError
|
||||
|
||||
async def close(self) -> None:
|
||||
pass
|
||||
44
tools/chat-bridge/chat_bridge/adapters/console.py
Normal file
44
tools/chat-bridge/chat_bridge/adapters/console.py
Normal file
@@ -0,0 +1,44 @@
|
||||
"""The smallest adapter: game chat is printed, and each line typed on stdin is posted in the game.
|
||||
|
||||
Lines are `Name: message`, or just `message` (sent as the configured default name). Useful for trying the bridge
|
||||
and as a template for new adapters.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import sys
|
||||
|
||||
from ..core import GameMessage, RemoteMessage, plain_line
|
||||
from .base import Adapter, Deliver
|
||||
|
||||
|
||||
def parse_line(line: str, default_name: str) -> RemoteMessage | None:
|
||||
line = line.strip()
|
||||
if not line:
|
||||
return None
|
||||
name, sep, text = line.partition(": ")
|
||||
if sep and 0 < len(name) <= 32 and text.strip():
|
||||
return RemoteMessage(author=name, text=text)
|
||||
return RemoteMessage(author=default_name, text=line)
|
||||
|
||||
|
||||
class ConsoleAdapter(Adapter):
|
||||
name = "console"
|
||||
default_label = "Console"
|
||||
|
||||
async def run(self, deliver: Deliver) -> None:
|
||||
if not self.config.get("read_stdin", True):
|
||||
return
|
||||
default_name = self.config.get("name", "Console")
|
||||
while True:
|
||||
line = await asyncio.to_thread(sys.stdin.readline)
|
||||
if not line: # end of input
|
||||
return
|
||||
msg = parse_line(line, default_name)
|
||||
if msg:
|
||||
await deliver(msg)
|
||||
|
||||
async def send(self, msg: GameMessage, skipped: int = 0) -> None:
|
||||
if skipped:
|
||||
print(f"({skipped} message(s) skipped: rate limit)", flush=True)
|
||||
print(plain_line(msg), flush=True)
|
||||
220
tools/chat-bridge/chat_bridge/adapters/discord.py
Normal file
220
tools/chat-bridge/chat_bridge/adapters/discord.py
Normal file
@@ -0,0 +1,220 @@
|
||||
"""Discord: one bot, one text channel.
|
||||
|
||||
Receiving uses the Discord gateway (a WebSocket) with the Guild Messages and Message Content intents; the Message
|
||||
Content intent must be switched on for the bot in the Discord developer portal. Sending uses a channel webhook when
|
||||
`webhook_url` is set (each player shows under their own name) or else the bot itself.
|
||||
|
||||
Only `websockets` is needed; the formatting functions at the top are plain Python and unit tested.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import random
|
||||
import re
|
||||
import urllib.error
|
||||
import urllib.request
|
||||
|
||||
from ..core import GameMessage, RemoteMessage, where
|
||||
from .base import Adapter, Deliver
|
||||
|
||||
log = logging.getLogger("chat_bridge.discord")
|
||||
|
||||
API = "https://discord.com/api/v10"
|
||||
GATEWAY = "wss://gateway.discord.gg/?v=10&encoding=json"
|
||||
INTENTS = (1 << 9) | (1 << 15) # GUILD_MESSAGES | MESSAGE_CONTENT
|
||||
FATAL_CLOSE_CODES = {4004: "the bot token is wrong", 4013: "invalid intents",
|
||||
4014: "the Message Content intent isn't enabled for this bot in the developer portal"}
|
||||
USER_AGENT = "DiscordBot (https://github.com/DarkflameUniverse/DarkflameServer, 1.0) dlu-chat-bridge"
|
||||
|
||||
# Nobody gets pinged by something said in the game
|
||||
NO_MENTIONS = {"parse": []}
|
||||
|
||||
_MARKDOWN = re.compile(r"([\\*_~`|>#\[\]])")
|
||||
_USER_MENTION = re.compile(r"<@!?(\d+)>")
|
||||
_ROLE_MENTION = re.compile(r"<@&\d+>")
|
||||
_CHANNEL_MENTION = re.compile(r"<#\d+>")
|
||||
_EMOJI = re.compile(r"<a?:(\w+):\d+>")
|
||||
_TIMESTAMP = re.compile(r"<t:(\d+)(?::\w)?>")
|
||||
|
||||
|
||||
# --- formatting (pure) ---
|
||||
|
||||
def escape_markdown(text: str) -> str:
|
||||
return _MARKDOWN.sub(r"\\\1", text)
|
||||
|
||||
|
||||
def webhook_username(name: str) -> str:
|
||||
"""Discord refuses webhook names containing "discord" or "clyde", and they are 1-80 characters."""
|
||||
name = re.sub(r"(?i)discord", "D1scord", name)
|
||||
name = re.sub(r"(?i)clyde", "Clyd3", name)
|
||||
name = name.strip() or "Player"
|
||||
return name[:80]
|
||||
|
||||
|
||||
def to_discord(msg: GameMessage, via_webhook: bool, skipped: int = 0) -> dict:
|
||||
"""The JSON body that posts one game message in Discord."""
|
||||
place = where(msg)
|
||||
note = f"-# {skipped} message(s) skipped (rate limit)\n" if skipped else ""
|
||||
if via_webhook:
|
||||
username = msg.sender_name + (f" ({place})" if place else "")
|
||||
return {"username": webhook_username(username), "content": note + escape_markdown(msg.message),
|
||||
"allowed_mentions": NO_MENTIONS}
|
||||
head = f"**{escape_markdown(msg.sender_name)}**" + (f" ({escape_markdown(place)})" if place else "")
|
||||
return {"content": f"{note}{head}: {escape_markdown(msg.message)}", "allowed_mentions": NO_MENTIONS}
|
||||
|
||||
|
||||
def display_name(event: dict) -> str:
|
||||
member = event.get("member") or {}
|
||||
author = event.get("author") or {}
|
||||
return member.get("nick") or author.get("global_name") or author.get("username") or "someone"
|
||||
|
||||
|
||||
def from_discord(event: dict) -> str:
|
||||
"""MESSAGE_CREATE content as plain text players can read: mentions and custom emoji spelled out."""
|
||||
names = {m.get("id"): (m.get("member") or {}).get("nick") or m.get("global_name") or m.get("username")
|
||||
for m in event.get("mentions", [])}
|
||||
text = event.get("content", "")
|
||||
text = _USER_MENTION.sub(lambda m: "@" + (names.get(m.group(1)) or "someone"), text)
|
||||
text = _ROLE_MENTION.sub("@role", text)
|
||||
text = _CHANNEL_MENTION.sub("#channel", text)
|
||||
text = _EMOJI.sub(r":\1:", text)
|
||||
text = _TIMESTAMP.sub("(a time)", text)
|
||||
extras = []
|
||||
if event.get("attachments"):
|
||||
extras.append("[attachment]")
|
||||
if event.get("sticker_items"):
|
||||
extras.append("[sticker]")
|
||||
return " ".join([text.strip(), *extras]).strip()
|
||||
|
||||
|
||||
def should_forward(event: dict, channel_id: str, own_user_id: str | None) -> bool:
|
||||
"""Whether a MESSAGE_CREATE goes into the game: our channel, a person (not a bot or webhook, so no loops)."""
|
||||
if str(event.get("channel_id")) != str(channel_id):
|
||||
return False
|
||||
if event.get("webhook_id"):
|
||||
return False
|
||||
author = event.get("author") or {}
|
||||
if author.get("bot") or (own_user_id and author.get("id") == own_user_id):
|
||||
return False
|
||||
return event.get("type", 0) in (0, 19) # a normal message or a reply
|
||||
|
||||
|
||||
# --- the adapter ---
|
||||
|
||||
class DiscordAdapter(Adapter):
|
||||
name = "discord"
|
||||
default_label = "Discord"
|
||||
|
||||
def __init__(self, config: dict):
|
||||
super().__init__(config)
|
||||
self.token = config.get("bot_token", "")
|
||||
self.channel_id = str(config.get("channel_id", ""))
|
||||
self.webhook_url = config.get("webhook_url", "")
|
||||
self.user_id: str | None = None
|
||||
if not self.webhook_url and not (self.token and self.channel_id):
|
||||
raise ValueError("discord needs bot_token and channel_id (to read and send) or webhook_url (to send only)")
|
||||
|
||||
# Sending
|
||||
|
||||
def _post(self, url: str, body: dict, bot_auth: bool) -> float:
|
||||
"""POST, returning seconds to wait and retry when Discord rate limits (0 when sent)."""
|
||||
headers = {"Content-Type": "application/json", "User-Agent": USER_AGENT}
|
||||
if bot_auth:
|
||||
headers["Authorization"] = f"Bot {self.token}"
|
||||
req = urllib.request.Request(url, data=json.dumps(body).encode(), method="POST", headers=headers)
|
||||
try:
|
||||
with urllib.request.urlopen(req, timeout=15) as resp:
|
||||
resp.read()
|
||||
return 0.0
|
||||
except urllib.error.HTTPError as e:
|
||||
if e.code == 429:
|
||||
try:
|
||||
return float(json.loads(e.read()).get("retry_after", 1.0))
|
||||
except ValueError:
|
||||
return 1.0
|
||||
raise
|
||||
|
||||
async def send(self, msg: GameMessage, skipped: int = 0) -> None:
|
||||
via_webhook = bool(self.webhook_url)
|
||||
body = to_discord(msg, via_webhook, skipped)
|
||||
url = self.webhook_url if via_webhook else f"{API}/channels/{self.channel_id}/messages"
|
||||
for _ in range(3):
|
||||
wait = await asyncio.to_thread(self._post, url, body, not via_webhook)
|
||||
if not wait:
|
||||
return
|
||||
await asyncio.sleep(wait)
|
||||
log.warning("gave up sending message %s to Discord: rate limited", msg.id)
|
||||
|
||||
# Receiving (the gateway)
|
||||
|
||||
async def run(self, deliver: Deliver) -> None:
|
||||
if not self.token or not self.channel_id:
|
||||
log.info("no bot_token/channel_id: sending to Discord only")
|
||||
return
|
||||
import websockets
|
||||
|
||||
session_id = resume_url = None
|
||||
seq = None
|
||||
backoff = 1.0
|
||||
while True:
|
||||
url = resume_url.rstrip("/") + "/?v=10&encoding=json" if session_id and resume_url else GATEWAY
|
||||
try:
|
||||
async with websockets.connect(url, max_size=2 ** 22) as ws:
|
||||
hello = json.loads(await ws.recv())
|
||||
interval = hello["d"]["heartbeat_interval"] / 1000.0
|
||||
if session_id:
|
||||
await ws.send(json.dumps({"op": 6, "d": {"token": self.token, "session_id": session_id, "seq": seq}}))
|
||||
else:
|
||||
await ws.send(json.dumps({"op": 2, "d": {"token": self.token, "intents": INTENTS,
|
||||
"properties": {"os": "linux", "browser": "dlu-chat-bridge", "device": "dlu-chat-bridge"}}}))
|
||||
last_seq = [seq]
|
||||
|
||||
async def heartbeat():
|
||||
await asyncio.sleep(interval * random.random())
|
||||
while True:
|
||||
await ws.send(json.dumps({"op": 1, "d": last_seq[0]}))
|
||||
await asyncio.sleep(interval)
|
||||
|
||||
beat = asyncio.create_task(heartbeat())
|
||||
try:
|
||||
async for raw in ws:
|
||||
event = json.loads(raw)
|
||||
op = event.get("op")
|
||||
if event.get("s") is not None:
|
||||
seq = last_seq[0] = event["s"]
|
||||
if op == 0:
|
||||
kind, data = event.get("t"), event.get("d") or {}
|
||||
if kind == "READY":
|
||||
session_id, resume_url = data["session_id"], data.get("resume_gateway_url")
|
||||
self.user_id = data["user"]["id"]
|
||||
backoff = 1.0
|
||||
log.info("connected to Discord as %s", data["user"].get("username"))
|
||||
elif kind == "RESUMED":
|
||||
backoff = 1.0
|
||||
elif kind == "MESSAGE_CREATE" and should_forward(data, self.channel_id, self.user_id):
|
||||
text = from_discord(data)
|
||||
if text:
|
||||
await deliver(RemoteMessage(author=display_name(data), text=text))
|
||||
elif op == 1:
|
||||
await ws.send(json.dumps({"op": 1, "d": last_seq[0]}))
|
||||
elif op == 7: # reconnect and resume
|
||||
break
|
||||
elif op == 9: # invalid session: start a new one unless it's resumable
|
||||
if not event.get("d"):
|
||||
session_id = resume_url = seq = None
|
||||
await asyncio.sleep(1 + random.random() * 4)
|
||||
break
|
||||
finally:
|
||||
beat.cancel()
|
||||
code = ws.close_code
|
||||
if code in FATAL_CLOSE_CODES:
|
||||
raise RuntimeError(f"Discord closed the connection: {FATAL_CLOSE_CODES[code]}")
|
||||
except (OSError, asyncio.TimeoutError, websockets.exceptions.WebSocketException) as e:
|
||||
code = getattr(e, "rcvd", None) and e.rcvd.code
|
||||
if code in FATAL_CLOSE_CODES:
|
||||
raise RuntimeError(f"Discord closed the connection: {FATAL_CLOSE_CODES[code]}") from None
|
||||
log.warning("Discord gateway: %s; reconnecting in %.0fs", e, backoff)
|
||||
await asyncio.sleep(backoff)
|
||||
backoff = min(backoff * 2, 60.0)
|
||||
39
tools/chat-bridge/chat_bridge/adapters/webhook.py
Normal file
39
tools/chat-bridge/chat_bridge/adapters/webhook.py
Normal file
@@ -0,0 +1,39 @@
|
||||
"""A send-only adapter: POSTs each game chat message as JSON to any URL (a Slack/Mattermost-style incoming webhook,
|
||||
an automation tool, your own service). It shows how little an adapter needs.
|
||||
|
||||
Body: {"text": "[Nimbus Station] Bob: hi", "sender": "Bob", "message": "hi", "channel": "zone", "zone_id": 1200,
|
||||
"zone_name": "Nimbus Station", "instance_id": 2, "id": 42}
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import urllib.request
|
||||
|
||||
from ..core import GameMessage, plain_line
|
||||
from .base import Adapter
|
||||
|
||||
log = logging.getLogger("chat_bridge.webhook")
|
||||
|
||||
|
||||
def webhook_body(msg: GameMessage, skipped: int = 0) -> dict:
|
||||
body = {"text": plain_line(msg), "sender": msg.sender_name, "message": msg.message, "channel": msg.channel,
|
||||
"zone_id": msg.zone_id, "zone_name": msg.zone_name, "instance_id": msg.instance_id, "id": msg.id}
|
||||
if skipped:
|
||||
body["skipped"] = skipped
|
||||
return body
|
||||
|
||||
|
||||
class WebhookAdapter(Adapter):
|
||||
name = "webhook"
|
||||
default_label = "Webhook"
|
||||
|
||||
def _post(self, body: dict) -> None:
|
||||
req = urllib.request.Request(self.config["url"], data=json.dumps(body).encode(), method="POST",
|
||||
headers={"Content-Type": "application/json", **self.config.get("headers", {})})
|
||||
with urllib.request.urlopen(req, timeout=15) as resp:
|
||||
resp.read()
|
||||
|
||||
async def send(self, msg: GameMessage, skipped: int = 0) -> None:
|
||||
await asyncio.to_thread(self._post, webhook_body(msg, skipped))
|
||||
186
tools/chat-bridge/chat_bridge/bridge.py
Normal file
186
tools/chat-bridge/chat_bridge/bridge.py
Normal file
@@ -0,0 +1,186 @@
|
||||
"""Runs the links: reads game chat from the dashboard (WebSocket to know when, GET /api/chat?after= to read), hands
|
||||
each message to the links that want it, and posts what the other side says with POST /api/chat/send."""
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
|
||||
from .adapters import Adapter
|
||||
from .config import Config, LinkConfig
|
||||
from .core import Cursor, GameMessage, RateLimiter, RemoteMessage, send_body, should_relay
|
||||
from .dashboard import ApiError, Dashboard
|
||||
|
||||
log = logging.getLogger("chat_bridge")
|
||||
|
||||
PAGE = 500 # the API's largest page
|
||||
QUEUE = 100 # messages waiting to go out per link before new ones are dropped
|
||||
|
||||
|
||||
def load_state(path: str) -> int:
|
||||
try:
|
||||
with open(path, encoding="utf-8") as f:
|
||||
return int(json.load(f).get("last_id", 0))
|
||||
except (OSError, ValueError):
|
||||
return 0
|
||||
|
||||
|
||||
def save_state(path: str, last_id: int) -> None:
|
||||
tmp = path + ".tmp"
|
||||
with open(tmp, "w", encoding="utf-8") as f:
|
||||
json.dump({"last_id": last_id}, f)
|
||||
os.replace(tmp, path)
|
||||
|
||||
|
||||
class Link:
|
||||
def __init__(self, cfg: LinkConfig, adapter: Adapter):
|
||||
self.cfg = cfg
|
||||
self.adapter = adapter
|
||||
self.out_limit = RateLimiter(cfg.to_remote_per_minute)
|
||||
self.in_limit = RateLimiter(cfg.to_game_per_minute)
|
||||
self.queue: asyncio.Queue = asyncio.Queue(QUEUE)
|
||||
|
||||
|
||||
class Bridge:
|
||||
def __init__(self, config: Config, dashboard: Dashboard, links: list[Link]):
|
||||
self.config = config
|
||||
self.dashboard = dashboard
|
||||
self.links = links
|
||||
self.cursor = Cursor(load_state(config.state_file))
|
||||
self.account_id = 0
|
||||
self.wake = asyncio.Event()
|
||||
self.ws_up = False
|
||||
# Ask the server for just one channel when every link wants the same single one (e.g. only zone chat), so
|
||||
# rows the bridge would drop anyway (whispers, for an account allowed to read them) never leave the server
|
||||
channels = {link.cfg.game.server_channel() for link in links}
|
||||
self.server_channel = channels.pop() if len(channels) == 1 else ""
|
||||
|
||||
# --- game -> other side ---
|
||||
|
||||
async def fetch(self) -> None:
|
||||
while True:
|
||||
page = await self.dashboard.chat_page(self.cursor.page_params(PAGE, self.server_channel))
|
||||
rows = page.get("messages", [])
|
||||
for msg in self.cursor.take(rows):
|
||||
self.dispatch(msg)
|
||||
# Rows the server filtered out still move the cursor on
|
||||
self.cursor.last_id = max(self.cursor.last_id, int(page.get("last_id", 0)))
|
||||
save_state(self.config.state_file, self.cursor.last_id)
|
||||
if len(rows) < PAGE:
|
||||
return
|
||||
|
||||
def dispatch(self, msg: GameMessage) -> None:
|
||||
for link in self.links:
|
||||
if not should_relay(msg, link.cfg.game, link.cfg.label, self.account_id):
|
||||
continue
|
||||
if not link.out_limit.allow():
|
||||
continue
|
||||
try:
|
||||
link.queue.put_nowait((msg, link.out_limit.take_dropped()))
|
||||
except asyncio.QueueFull:
|
||||
link.out_limit.dropped += 1
|
||||
|
||||
async def send_loop(self, link: Link) -> None:
|
||||
while True:
|
||||
msg, skipped = await link.queue.get()
|
||||
try:
|
||||
await link.adapter.send(msg, skipped)
|
||||
except Exception as e: # one failed send mustn't stop the link
|
||||
log.warning("%s: couldn't send message %s: %s", link.cfg.label, msg.id, e)
|
||||
|
||||
async def read_loop(self) -> None:
|
||||
backoff = 1.0
|
||||
while True:
|
||||
try:
|
||||
await self.fetch()
|
||||
backoff = 1.0
|
||||
except ApiError as e:
|
||||
if e.status in (401, 403):
|
||||
raise RuntimeError(f"the dashboard refused the token reading chat ({e}); it needs chat_view and api_access") from None
|
||||
log.warning("reading chat failed: %s", e)
|
||||
except OSError as e:
|
||||
log.warning("reading chat failed: %s; retrying in %.0fs", e, backoff)
|
||||
await asyncio.sleep(backoff)
|
||||
backoff = min(backoff * 2, 60.0)
|
||||
continue
|
||||
# The WebSocket says when there's something new; without it, poll. Either way re-read now and then.
|
||||
timeout = self.config.resync_seconds if self.ws_up else self.config.poll_seconds
|
||||
try:
|
||||
await asyncio.wait_for(self.wake.wait(), timeout)
|
||||
except asyncio.TimeoutError:
|
||||
pass
|
||||
self.wake.clear()
|
||||
|
||||
async def ws_loop(self) -> None:
|
||||
backoff = 1.0
|
||||
|
||||
def connected():
|
||||
nonlocal backoff
|
||||
self.ws_up, backoff = True, 1.0
|
||||
log.info("live: subscribed to chat_message")
|
||||
self.wake.set() # catch up on anything said while the socket was down
|
||||
|
||||
while True:
|
||||
try:
|
||||
await self.dashboard.watch_chat(lambda _event: self.wake.set(), connected)
|
||||
log.warning("WebSocket closed; polling every %.0fs until it's back", self.config.poll_seconds)
|
||||
except ImportError:
|
||||
log.warning("the websockets package isn't installed: polling every %.0fs", self.config.poll_seconds)
|
||||
return
|
||||
except ApiError as e:
|
||||
log.warning("%s; polling instead", e)
|
||||
return
|
||||
except Exception as e:
|
||||
log.warning("WebSocket: %s; polling every %.0fs, reconnecting in %.0fs", e, self.config.poll_seconds, backoff)
|
||||
self.ws_up = False
|
||||
self.wake.set()
|
||||
await asyncio.sleep(backoff)
|
||||
backoff = min(backoff * 2, 60.0)
|
||||
|
||||
# --- other side -> game ---
|
||||
|
||||
def deliverer(self, link: Link):
|
||||
async def deliver(remote: RemoteMessage) -> None:
|
||||
if not link.in_limit.allow():
|
||||
log.info("%s: rate limit, not posting %s's message in the game", link.cfg.label, remote.author)
|
||||
return
|
||||
body = send_body(remote, link.cfg.label, link.cfg.game)
|
||||
if body is None:
|
||||
return
|
||||
try:
|
||||
await self.dashboard.send_chat(body)
|
||||
except (ApiError, OSError) as e:
|
||||
log.warning("%s: couldn't post in the game: %s", link.cfg.label, e)
|
||||
return deliver
|
||||
|
||||
async def run_adapter(self, link: Link) -> None:
|
||||
await link.adapter.run(self.deliverer(link))
|
||||
# An adapter that only sends (or reached the end of its input) keeps relaying out
|
||||
await asyncio.Event().wait()
|
||||
|
||||
# --- start ---
|
||||
|
||||
async def start(self) -> None:
|
||||
me = await self.dashboard.me()
|
||||
if not me.get("valid"):
|
||||
raise RuntimeError("the API token isn't valid (expired, revoked, or the account was signed out everywhere)")
|
||||
self.account_id = int(me.get("accountId", 0))
|
||||
log.info("signed in as %s (GM %s)", me.get("username"), me.get("gmLevel"))
|
||||
if not self.cursor.started:
|
||||
# First run: start from now rather than replaying the whole log
|
||||
page = await self.dashboard.chat_page({"limit": "1"})
|
||||
self.cursor.last_id = int(page.get("last_id", 0))
|
||||
save_state(self.config.state_file, self.cursor.last_id)
|
||||
log.info("relaying chat after id %d", self.cursor.last_id)
|
||||
|
||||
async def run(self) -> None:
|
||||
await self.start()
|
||||
tasks = [self.read_loop(), self.ws_loop()]
|
||||
for link in self.links:
|
||||
tasks += [self.send_loop(link), self.run_adapter(link)]
|
||||
try:
|
||||
await asyncio.gather(*tasks)
|
||||
finally:
|
||||
for link in self.links:
|
||||
await link.adapter.close()
|
||||
77
tools/chat-bridge/chat_bridge/config.py
Normal file
77
tools/chat-bridge/chat_bridge/config.py
Normal file
@@ -0,0 +1,77 @@
|
||||
"""Reading the bridge's JSON config. Any string value "$NAME" or "${NAME}" is read from the environment, so secrets
|
||||
can stay out of the file; DLU_DASHBOARD_URL and DLU_API_TOKEN override the dashboard settings outright."""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
from dataclasses import dataclass, field
|
||||
|
||||
from .core import LABEL_PATTERN, GameFilter
|
||||
|
||||
_ENV = re.compile(r"^\$\{?([A-Za-z_][A-Za-z0-9_]*)\}?$")
|
||||
|
||||
|
||||
def expand_env(value, environ=None):
|
||||
environ = os.environ if environ is None else environ
|
||||
if isinstance(value, dict):
|
||||
return {k: expand_env(v, environ) for k, v in value.items()}
|
||||
if isinstance(value, list):
|
||||
return [expand_env(v, environ) for v in value]
|
||||
if isinstance(value, str):
|
||||
m = _ENV.match(value)
|
||||
if m:
|
||||
return environ.get(m.group(1), "")
|
||||
return value
|
||||
|
||||
|
||||
@dataclass
|
||||
class LinkConfig:
|
||||
adapter: str
|
||||
label: str
|
||||
game: GameFilter
|
||||
settings: dict
|
||||
to_game_per_minute: float = 20
|
||||
to_remote_per_minute: float = 60
|
||||
|
||||
|
||||
@dataclass
|
||||
class Config:
|
||||
url: str
|
||||
token: str
|
||||
links: list = field(default_factory=list)
|
||||
state_file: str = "chat-bridge.state.json"
|
||||
poll_seconds: float = 5
|
||||
resync_seconds: float = 60
|
||||
|
||||
|
||||
def parse_config(raw: dict, environ=None, default_labels: dict | None = None) -> Config:
|
||||
environ = os.environ if environ is None else environ
|
||||
raw = expand_env(raw, environ)
|
||||
dash = raw.get("dashboard", {})
|
||||
url = environ.get("DLU_DASHBOARD_URL") or dash.get("url") or "http://localhost:2006"
|
||||
token = environ.get("DLU_API_TOKEN") or dash.get("token", "")
|
||||
if not token:
|
||||
raise ValueError("no API token: set dashboard.token in the config or DLU_API_TOKEN")
|
||||
links = []
|
||||
for i, link in enumerate(raw.get("links", [])):
|
||||
adapter = link.get("adapter")
|
||||
if not adapter:
|
||||
raise ValueError(f"links[{i}] has no adapter")
|
||||
label = link.get("label") or (default_labels or {}).get(adapter, adapter.capitalize())
|
||||
if not LABEL_PATTERN.match(label):
|
||||
raise ValueError(f"links[{i}].label: up to 16 letters, digits, spaces, - or _")
|
||||
rate = link.get("rate", {})
|
||||
links.append(LinkConfig(adapter=adapter, label=label, game=GameFilter.from_config(link.get("game", {})),
|
||||
settings=link.get(adapter, {}),
|
||||
to_game_per_minute=float(rate.get("to_game_per_minute", 20)),
|
||||
to_remote_per_minute=float(rate.get("to_remote_per_minute", 60))))
|
||||
if not links:
|
||||
raise ValueError("no links configured")
|
||||
return Config(url=url, token=token, links=links, state_file=raw.get("state_file", "chat-bridge.state.json"),
|
||||
poll_seconds=float(raw.get("poll_seconds", 5)), resync_seconds=float(raw.get("resync_seconds", 60)))
|
||||
|
||||
|
||||
def load_config(path: str, default_labels: dict | None = None) -> Config:
|
||||
with open(path, encoding="utf-8") as f:
|
||||
return parse_config(json.load(f), default_labels=default_labels)
|
||||
222
tools/chat-bridge/chat_bridge/core.py
Normal file
222
tools/chat-bridge/chat_bridge/core.py
Normal file
@@ -0,0 +1,222 @@
|
||||
"""The parts of the bridge that don't touch the network: what to relay, how to word it, where to resume.
|
||||
|
||||
Everything here is plain Python so it can be unit tested without a dashboard or a chat service.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import re
|
||||
import time
|
||||
from dataclasses import dataclass, field
|
||||
from typing import Iterable, Optional
|
||||
|
||||
# Limits of POST /api/chat/send (dDashboardServer/routes/ChatRoutes.cpp)
|
||||
GAME_MESSAGE_MAX = 300
|
||||
GAME_NAME_MAX = 32
|
||||
LABEL_PATTERN = re.compile(r"^[A-Za-z0-9 _-]{1,16}$")
|
||||
|
||||
PRIVATE_CHANNELS = frozenset({"whisper", "team"})
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class GameMessage:
|
||||
"""One row of GET /api/chat or one chat_message WebSocket event."""
|
||||
id: int
|
||||
channel: str
|
||||
sender_name: str
|
||||
message: str
|
||||
zone_id: int = 0
|
||||
zone_name: str = ""
|
||||
instance_id: int = 0
|
||||
account_id: int = 0
|
||||
blocked: bool = False
|
||||
time: int = 0
|
||||
|
||||
@classmethod
|
||||
def from_json(cls, data: dict) -> "GameMessage":
|
||||
return cls(
|
||||
id=int(data.get("id", 0)),
|
||||
channel=str(data.get("channel", "")),
|
||||
sender_name=str(data.get("sender_name", "")),
|
||||
message=str(data.get("message", "")),
|
||||
zone_id=int(data.get("zone_id", 0) or 0),
|
||||
zone_name=str(data.get("zone_name", "") or ""),
|
||||
instance_id=int(data.get("instance_id", 0) or 0),
|
||||
account_id=int(data.get("account_id", 0) or 0),
|
||||
blocked=bool(data.get("blocked", False)),
|
||||
time=int(data.get("time", 0) or 0),
|
||||
)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class RemoteMessage:
|
||||
"""A message from the other side (Discord, the console, ...) to post in the game."""
|
||||
author: str
|
||||
text: str
|
||||
|
||||
|
||||
@dataclass
|
||||
class GameFilter:
|
||||
"""Which game chat a link relays out, and where messages coming in are posted."""
|
||||
channels: frozenset = frozenset({"zone"})
|
||||
zones: frozenset = frozenset() # empty: every zone
|
||||
instances: frozenset = frozenset() # empty: every world of those zones
|
||||
allow_private: bool = False # whispers and team chat; off unless explicitly turned on
|
||||
send_zone: int = 0 # 0: every zone
|
||||
send_instance: Optional[int] = None # with send_zone: only that world
|
||||
|
||||
@classmethod
|
||||
def from_config(cls, cfg: dict) -> "GameFilter":
|
||||
channels = frozenset(cfg.get("channels", ["zone"]))
|
||||
allow_private = bool(cfg.get("allow_private", False))
|
||||
private = channels & PRIVATE_CHANNELS
|
||||
if private and not allow_private:
|
||||
raise ValueError(f"channels {sorted(private)} are private: set \"allow_private\": true to relay them")
|
||||
send_instance = cfg.get("send_instance")
|
||||
return cls(
|
||||
channels=channels,
|
||||
zones=frozenset(int(z) for z in cfg.get("zones", [])),
|
||||
instances=frozenset(int(i) for i in cfg.get("instances", [])),
|
||||
allow_private=allow_private,
|
||||
send_zone=int(cfg.get("send_zone", 0) or 0),
|
||||
send_instance=int(send_instance) if send_instance is not None else None,
|
||||
)
|
||||
|
||||
def server_channel(self) -> str:
|
||||
"""The channel= to ask the API for, so rows the bridge would drop (like whispers) never leave the server."""
|
||||
return next(iter(self.channels)) if len(self.channels) == 1 else ""
|
||||
|
||||
|
||||
def is_own_message(msg: GameMessage, label: str, own_account_id: int = 0) -> bool:
|
||||
"""A message this bridge (or anything posting with the same label) sent into the game.
|
||||
|
||||
POST /api/chat/send logs what it posts in the "web" channel as "[label] name", from the token's account.
|
||||
"""
|
||||
if msg.channel != "web":
|
||||
return False
|
||||
if msg.sender_name.startswith(f"[{label}] "):
|
||||
return True
|
||||
return bool(own_account_id) and msg.account_id == own_account_id
|
||||
|
||||
|
||||
def should_relay(msg: GameMessage, flt: GameFilter, label: str, own_account_id: int = 0) -> bool:
|
||||
"""Whether a game message goes to the other side of a link."""
|
||||
if msg.blocked: # the chat filter stopped it: no player saw it, so neither should anyone else
|
||||
return False
|
||||
if msg.channel in PRIVATE_CHANNELS and not flt.allow_private:
|
||||
return False
|
||||
if msg.channel not in flt.channels:
|
||||
return False
|
||||
if is_own_message(msg, label, own_account_id): # loop prevention
|
||||
return False
|
||||
if flt.zones and msg.zone_id not in flt.zones:
|
||||
return False
|
||||
if flt.instances and msg.instance_id not in flt.instances:
|
||||
return False
|
||||
return bool(msg.message.strip())
|
||||
|
||||
|
||||
def one_line(text: str) -> str:
|
||||
return re.sub(r"\s+", " ", text).strip()
|
||||
|
||||
|
||||
def truncate(text: str, limit: int) -> str:
|
||||
return text if len(text) <= limit else text[: limit - 3].rstrip() + "..."
|
||||
|
||||
|
||||
def game_name(author: str, fallback: str = "someone") -> str:
|
||||
"""A display name the game endpoint accepts: one line, no brackets, 1-32 characters."""
|
||||
name = one_line(author).replace("[", "(").replace("]", ")")
|
||||
return truncate(name, GAME_NAME_MAX) or fallback
|
||||
|
||||
|
||||
def game_text(text: str) -> str:
|
||||
"""Message text the game endpoint accepts: one line, 1-300 characters (empty: nothing to send)."""
|
||||
return truncate(one_line(text), GAME_MESSAGE_MAX)
|
||||
|
||||
|
||||
def send_body(remote: RemoteMessage, label: str, flt: GameFilter) -> Optional[dict]:
|
||||
"""The JSON for POST /api/chat/send, or None when there's nothing to post."""
|
||||
if not LABEL_PATTERN.match(label):
|
||||
raise ValueError("label is up to 16 letters, digits, spaces, - or _")
|
||||
text = game_text(remote.text)
|
||||
if not text:
|
||||
return None
|
||||
body = {"message": text, "name": game_name(remote.author), "label": label}
|
||||
if flt.send_zone:
|
||||
body["zone"] = flt.send_zone
|
||||
if flt.send_instance is not None:
|
||||
body["instance"] = flt.send_instance
|
||||
return body
|
||||
|
||||
|
||||
def where(msg: GameMessage) -> str:
|
||||
"""Where in the game a message was said, for the other side."""
|
||||
if not msg.zone_id:
|
||||
return ""
|
||||
return msg.zone_name or f"zone {msg.zone_id}"
|
||||
|
||||
|
||||
def plain_line(msg: GameMessage) -> str:
|
||||
"""A plain one-line rendering: `[Nimbus Station] Bob: hi` (web messages already carry their [label])."""
|
||||
place = where(msg)
|
||||
prefix = f"[{place}] " if place else ""
|
||||
return f"{prefix}{msg.sender_name}: {msg.message}"
|
||||
|
||||
|
||||
class Cursor:
|
||||
"""Remembers the newest chat id handled so a restart or a dropped socket resumes without gaps or repeats.
|
||||
|
||||
The bridge always reads with GET /api/chat?after=<last_id> (the WebSocket only says when to read), so ids
|
||||
arrive in order; take() still drops anything at or below last_id in case pages overlap.
|
||||
"""
|
||||
|
||||
def __init__(self, last_id: int = 0):
|
||||
self.last_id = int(last_id)
|
||||
|
||||
@property
|
||||
def started(self) -> bool:
|
||||
return self.last_id > 0
|
||||
|
||||
def take(self, rows: Iterable[dict]) -> list[GameMessage]:
|
||||
fresh = sorted((GameMessage.from_json(r) for r in rows), key=lambda m: m.id)
|
||||
out = []
|
||||
for msg in fresh:
|
||||
if msg.id <= self.last_id:
|
||||
continue
|
||||
out.append(msg)
|
||||
self.last_id = msg.id
|
||||
return out
|
||||
|
||||
def page_params(self, limit: int, channel: str = "") -> dict:
|
||||
params = {"after": str(self.last_id), "limit": str(limit)}
|
||||
if channel:
|
||||
params["channel"] = channel
|
||||
return params
|
||||
|
||||
|
||||
@dataclass
|
||||
class RateLimiter:
|
||||
"""A token bucket: `per_minute` messages a minute, with bursts of up to `burst`."""
|
||||
per_minute: float
|
||||
burst: int = 5
|
||||
_tokens: float = field(default=-1.0, repr=False)
|
||||
_stamp: float = field(default=0.0, repr=False)
|
||||
dropped: int = 0
|
||||
|
||||
def allow(self, now: Optional[float] = None) -> bool:
|
||||
now = time.monotonic() if now is None else now
|
||||
if self.per_minute <= 0:
|
||||
return True
|
||||
if self._tokens < 0:
|
||||
self._tokens, self._stamp = float(self.burst), now
|
||||
self._tokens = min(float(self.burst), self._tokens + (now - self._stamp) * self.per_minute / 60.0)
|
||||
self._stamp = now
|
||||
if self._tokens >= 1.0:
|
||||
self._tokens -= 1.0
|
||||
return True
|
||||
self.dropped += 1
|
||||
return False
|
||||
|
||||
def take_dropped(self) -> int:
|
||||
n, self.dropped = self.dropped, 0
|
||||
return n
|
||||
93
tools/chat-bridge/chat_bridge/dashboard.py
Normal file
93
tools/chat-bridge/chat_bridge/dashboard.py
Normal file
@@ -0,0 +1,93 @@
|
||||
"""Talks to the dashboard's public API: GET /api/chat, POST /api/chat/send, and the /ws chat_message topic."""
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import urllib.error
|
||||
import urllib.parse
|
||||
import urllib.request
|
||||
|
||||
log = logging.getLogger("chat_bridge.dashboard")
|
||||
|
||||
|
||||
class ApiError(Exception):
|
||||
def __init__(self, status: int, message: str):
|
||||
super().__init__(f"HTTP {status}: {message}")
|
||||
self.status = status
|
||||
|
||||
|
||||
class Dashboard:
|
||||
def __init__(self, url: str, token: str, timeout: float = 15.0):
|
||||
self.url = url.rstrip("/")
|
||||
self.token = token
|
||||
self.timeout = timeout
|
||||
|
||||
# --- HTTP (stdlib, run in a thread so the event loop keeps going) ---
|
||||
|
||||
def _request(self, method: str, path: str, params: dict | None = None, body: dict | None = None) -> dict:
|
||||
url = self.url + path
|
||||
if params:
|
||||
url += "?" + urllib.parse.urlencode(params)
|
||||
data = json.dumps(body).encode() if body is not None else None
|
||||
req = urllib.request.Request(url, data=data, method=method, headers={
|
||||
"Authorization": f"Bearer {self.token}",
|
||||
"Accept": "application/json",
|
||||
**({"Content-Type": "application/json"} if data is not None else {}),
|
||||
})
|
||||
try:
|
||||
with urllib.request.urlopen(req, timeout=self.timeout) as resp:
|
||||
return json.loads(resp.read() or b"{}")
|
||||
except urllib.error.HTTPError as e:
|
||||
try:
|
||||
message = json.loads(e.read()).get("message", e.reason)
|
||||
except Exception:
|
||||
message = e.reason
|
||||
raise ApiError(e.code, str(message)) from None
|
||||
|
||||
async def get(self, path: str, params: dict | None = None) -> dict:
|
||||
return await asyncio.to_thread(self._request, "GET", path, params)
|
||||
|
||||
async def post(self, path: str, body: dict) -> dict:
|
||||
return await asyncio.to_thread(self._request, "POST", path, None, body)
|
||||
|
||||
async def me(self) -> dict:
|
||||
return await self.get("/api/auth/me")
|
||||
|
||||
async def chat_page(self, params: dict) -> dict:
|
||||
return await self.get("/api/chat", params)
|
||||
|
||||
async def send_chat(self, body: dict) -> dict:
|
||||
return await self.post("/api/chat/send", body)
|
||||
|
||||
# --- WebSocket ---
|
||||
|
||||
def ws_url(self) -> str:
|
||||
parts = urllib.parse.urlsplit(self.url)
|
||||
scheme = "wss" if parts.scheme == "https" else "ws"
|
||||
return urllib.parse.urlunsplit((scheme, parts.netloc, parts.path.rstrip("/") + "/ws", "", ""))
|
||||
|
||||
async def watch_chat(self, on_message, on_connected=None) -> None:
|
||||
"""Connect to /ws, subscribe to chat_message and call on_message(event) for each. Returns when the socket closes."""
|
||||
import websockets # only needed for live updates; polling works without it
|
||||
|
||||
headers = {"Authorization": f"Bearer {self.token}"}
|
||||
try:
|
||||
from websockets.asyncio.client import connect # websockets 13+
|
||||
ctx = connect(self.ws_url(), additional_headers=headers, open_timeout=self.timeout)
|
||||
except ImportError: # older websockets
|
||||
ctx = websockets.connect(self.ws_url(), extra_headers=headers, open_timeout=self.timeout)
|
||||
async with ctx as ws:
|
||||
await ws.send(json.dumps({"event": "subscribe", "subscription": "chat_message"}))
|
||||
reply = json.loads(await asyncio.wait_for(ws.recv(), self.timeout))
|
||||
if reply.get("error"):
|
||||
raise ApiError(403, f"can't subscribe to chat_message: {reply['error']} (the account needs chat_view)")
|
||||
if on_connected:
|
||||
on_connected()
|
||||
async for raw in ws:
|
||||
try:
|
||||
event = json.loads(raw)
|
||||
except ValueError:
|
||||
continue
|
||||
if event.get("event") == "chat_message":
|
||||
on_message(event)
|
||||
11
tools/chat-bridge/config.console.json
Normal file
11
tools/chat-bridge/config.console.json
Normal file
@@ -0,0 +1,11 @@
|
||||
{
|
||||
"dashboard": { "url": "http://localhost:2006", "token": "$DLU_API_TOKEN" },
|
||||
"state_file": "console.state.json",
|
||||
"links": [
|
||||
{
|
||||
"adapter": "console",
|
||||
"console": { "name": "Console" },
|
||||
"game": { "channels": ["zone", "web"] }
|
||||
}
|
||||
]
|
||||
}
|
||||
27
tools/chat-bridge/config.example.json
Normal file
27
tools/chat-bridge/config.example.json
Normal file
@@ -0,0 +1,27 @@
|
||||
{
|
||||
"dashboard": {
|
||||
"url": "http://localhost:2006",
|
||||
"token": "$DLU_API_TOKEN"
|
||||
},
|
||||
"state_file": "chat-bridge.state.json",
|
||||
"poll_seconds": 5,
|
||||
"resync_seconds": 60,
|
||||
"links": [
|
||||
{
|
||||
"adapter": "discord",
|
||||
"label": "Discord",
|
||||
"discord": {
|
||||
"bot_token": "$DISCORD_BOT_TOKEN",
|
||||
"channel_id": "123456789012345678",
|
||||
"webhook_url": "$DISCORD_WEBHOOK_URL"
|
||||
},
|
||||
"game": {
|
||||
"channels": ["zone"],
|
||||
"zones": [],
|
||||
"instances": [],
|
||||
"send_zone": 0
|
||||
},
|
||||
"rate": { "to_game_per_minute": 20, "to_remote_per_minute": 60 }
|
||||
}
|
||||
]
|
||||
}
|
||||
2
tools/chat-bridge/requirements.txt
Normal file
2
tools/chat-bridge/requirements.txt
Normal file
@@ -0,0 +1,2 @@
|
||||
# Only for live updates over the WebSocket (and the Discord adapter); without it the bridge polls.
|
||||
websockets>=10
|
||||
276
tools/chat-bridge/tests/test_core.py
Normal file
276
tools/chat-bridge/tests/test_core.py
Normal file
@@ -0,0 +1,276 @@
|
||||
import asyncio
|
||||
import os
|
||||
import sys
|
||||
import tempfile
|
||||
import unittest
|
||||
|
||||
sys.path.insert(0, os.path.join(os.path.dirname(__file__), ".."))
|
||||
|
||||
from chat_bridge.adapters.console import parse_line # noqa: E402
|
||||
from chat_bridge.adapters.webhook import webhook_body # noqa: E402
|
||||
from chat_bridge.bridge import Bridge, Link, load_state, save_state # noqa: E402
|
||||
from chat_bridge.config import Config, LinkConfig, expand_env, parse_config # noqa: E402
|
||||
from chat_bridge.core import (Cursor, GameFilter, GameMessage, RateLimiter, RemoteMessage, game_name, # noqa: E402
|
||||
is_own_message, plain_line, send_body, should_relay)
|
||||
|
||||
|
||||
def row(id, channel="zone", sender="Bob", message="hi", zone=1200, instance=2, account=5, blocked=False, zone_name="Nimbus Station"):
|
||||
return {"id": id, "channel": channel, "sender_name": sender, "message": message, "zone_id": zone, "zone_name": zone_name,
|
||||
"instance_id": instance, "account_id": account, "blocked": blocked, "time": 0}
|
||||
|
||||
|
||||
def msg(**kw):
|
||||
return GameMessage.from_json(row(kw.pop("id", 1), **kw))
|
||||
|
||||
|
||||
class RelayRules(unittest.TestCase):
|
||||
def test_zone_chat_by_default(self):
|
||||
f = GameFilter()
|
||||
self.assertTrue(should_relay(msg(), f, "Discord"))
|
||||
self.assertFalse(should_relay(msg(channel="web", sender="[Web] admin"), f, "Discord"))
|
||||
|
||||
def test_private_chat_never_without_opt_in(self):
|
||||
self.assertFalse(should_relay(msg(channel="whisper"), GameFilter(channels=frozenset({"zone", "whisper"})), "Discord"))
|
||||
self.assertFalse(should_relay(msg(channel="team"), GameFilter(channels=frozenset({"team"})), "Discord"))
|
||||
with self.assertRaises(ValueError):
|
||||
GameFilter.from_config({"channels": ["zone", "whisper"]})
|
||||
f = GameFilter.from_config({"channels": ["whisper"], "allow_private": True})
|
||||
self.assertTrue(should_relay(msg(channel="whisper"), f, "Discord"))
|
||||
|
||||
def test_blocked_messages_are_not_relayed(self):
|
||||
self.assertFalse(should_relay(msg(blocked=True), GameFilter(), "Discord"))
|
||||
|
||||
def test_zone_and_instance_filters(self):
|
||||
f = GameFilter(zones=frozenset({1100}))
|
||||
self.assertFalse(should_relay(msg(zone=1200), f, "Discord"))
|
||||
self.assertTrue(should_relay(msg(zone=1100), f, "Discord"))
|
||||
f = GameFilter(zones=frozenset({1200}), instances=frozenset({3}))
|
||||
self.assertFalse(should_relay(msg(zone=1200, instance=2), f, "Discord"))
|
||||
|
||||
def test_empty_messages_skipped(self):
|
||||
self.assertFalse(should_relay(msg(message=" "), GameFilter(), "Discord"))
|
||||
|
||||
def test_server_channel(self):
|
||||
self.assertEqual(GameFilter().server_channel(), "zone")
|
||||
self.assertEqual(GameFilter(channels=frozenset({"zone", "web"})).server_channel(), "")
|
||||
|
||||
|
||||
class LoopPrevention(unittest.TestCase):
|
||||
def test_own_label_is_ignored(self):
|
||||
f = GameFilter(channels=frozenset({"zone", "web"}))
|
||||
own = msg(channel="web", sender="[Discord] Alice", account=99)
|
||||
self.assertTrue(is_own_message(own, "Discord"))
|
||||
self.assertFalse(should_relay(own, f, "Discord"))
|
||||
|
||||
def test_own_account_is_ignored(self):
|
||||
f = GameFilter(channels=frozenset({"zone", "web"}))
|
||||
other_label = msg(channel="web", sender="[Matrix] Carol", account=99)
|
||||
self.assertFalse(should_relay(other_label, f, "Discord", own_account_id=99))
|
||||
self.assertTrue(should_relay(other_label, f, "Discord", own_account_id=7))
|
||||
|
||||
def test_player_named_like_a_label_is_not_own(self):
|
||||
# Only the web channel carries labels; a player can't produce one
|
||||
self.assertFalse(is_own_message(msg(channel="zone", sender="[Discord] x"), "Discord"))
|
||||
|
||||
|
||||
class Formatting(unittest.TestCase):
|
||||
def test_plain_line(self):
|
||||
self.assertEqual(plain_line(msg()), "[Nimbus Station] Bob: hi")
|
||||
self.assertEqual(plain_line(msg(zone=0, zone_name="")), "Bob: hi")
|
||||
self.assertEqual(plain_line(msg(zone=1234, zone_name="")), "[zone 1234] Bob: hi")
|
||||
|
||||
def test_send_body(self):
|
||||
body = send_body(RemoteMessage("Alice", "hello\nthere"), "Discord", GameFilter())
|
||||
self.assertEqual(body, {"message": "hello there", "name": "Alice", "label": "Discord"})
|
||||
|
||||
def test_send_body_targets(self):
|
||||
body = send_body(RemoteMessage("A", "x"), "Discord", GameFilter(send_zone=1200, send_instance=2))
|
||||
self.assertEqual((body["zone"], body["instance"]), (1200, 2))
|
||||
body = send_body(RemoteMessage("A", "x"), "Discord", GameFilter(send_instance=2))
|
||||
self.assertNotIn("zone", body)
|
||||
self.assertNotIn("instance", body)
|
||||
|
||||
def test_send_body_limits(self):
|
||||
body = send_body(RemoteMessage("[mod] " + "n" * 50, "x" * 400), "Discord", GameFilter())
|
||||
self.assertEqual(len(body["message"]), 300)
|
||||
self.assertTrue(body["message"].endswith("..."))
|
||||
self.assertLessEqual(len(body["name"]), 32)
|
||||
self.assertNotIn("[", body["name"])
|
||||
self.assertIsNone(send_body(RemoteMessage("A", " \n "), "Discord", GameFilter()))
|
||||
with self.assertRaises(ValueError):
|
||||
send_body(RemoteMessage("A", "x"), "Bad[label]", GameFilter())
|
||||
|
||||
def test_game_name_fallback(self):
|
||||
self.assertEqual(game_name(" \n"), "someone")
|
||||
|
||||
def test_console_parse(self):
|
||||
self.assertEqual(parse_line("Alice: hi there\n", "Console"), RemoteMessage("Alice", "hi there"))
|
||||
self.assertEqual(parse_line("just text\n", "Console"), RemoteMessage("Console", "just text"))
|
||||
self.assertIsNone(parse_line(" \n", "Console"))
|
||||
|
||||
def test_webhook_body(self):
|
||||
body = webhook_body(msg(id=42), skipped=3)
|
||||
self.assertEqual(body["text"], "[Nimbus Station] Bob: hi")
|
||||
self.assertEqual((body["id"], body["skipped"]), (42, 3))
|
||||
|
||||
|
||||
class Resume(unittest.TestCase):
|
||||
def test_take_in_order_and_advance(self):
|
||||
c = Cursor(10)
|
||||
out = c.take([row(12), row(11), row(13)])
|
||||
self.assertEqual([m.id for m in out], [11, 12, 13])
|
||||
self.assertEqual(c.last_id, 13)
|
||||
|
||||
def test_no_duplicates_across_overlapping_pages(self):
|
||||
c = Cursor(10)
|
||||
c.take([row(11), row(12)])
|
||||
out = c.take([row(12), row(13)])
|
||||
self.assertEqual([m.id for m in out], [13])
|
||||
|
||||
def test_page_params(self):
|
||||
c = Cursor(7)
|
||||
self.assertEqual(c.page_params(500, "zone"), {"after": "7", "limit": "500", "channel": "zone"})
|
||||
self.assertEqual(c.page_params(500), {"after": "7", "limit": "500"})
|
||||
|
||||
def test_state_file_roundtrip(self):
|
||||
with tempfile.TemporaryDirectory() as d:
|
||||
path = os.path.join(d, "s.json")
|
||||
self.assertEqual(load_state(path), 0)
|
||||
save_state(path, 1234)
|
||||
self.assertEqual(load_state(path), 1234)
|
||||
|
||||
|
||||
class FakeDashboard:
|
||||
"""Serves GET /api/chat from a list, like the server: after=, limit=, channel=."""
|
||||
|
||||
def __init__(self, rows):
|
||||
self.rows = rows
|
||||
self.sent = []
|
||||
|
||||
async def chat_page(self, params):
|
||||
after, limit = int(params["after"]), int(params["limit"])
|
||||
rows = [r for r in self.rows if r["id"] > after and (not params.get("channel") or r["channel"] == params["channel"])][:limit]
|
||||
return {"messages": rows, "last_id": max([after] + [r["id"] for r in rows])}
|
||||
|
||||
async def send_chat(self, body):
|
||||
self.sent.append(body)
|
||||
return {"requestId": 1}
|
||||
|
||||
|
||||
class Recorder:
|
||||
name = "rec"
|
||||
|
||||
def __init__(self):
|
||||
self.got = []
|
||||
|
||||
async def send(self, m, skipped=0):
|
||||
self.got.append(m)
|
||||
|
||||
|
||||
class BridgeFlow(unittest.TestCase):
|
||||
def make(self, rows, state, game=None, per_min=0):
|
||||
cfg = Config(url="x", token="t", state_file=state)
|
||||
link = Link(LinkConfig("rec", "Discord", game or GameFilter(), {}, per_min, per_min), Recorder())
|
||||
dash = FakeDashboard(rows)
|
||||
return Bridge(cfg, dash, [link]), link, dash
|
||||
|
||||
def test_resume_after_restart(self):
|
||||
rows = [row(i) for i in range(1, 8)]
|
||||
with tempfile.TemporaryDirectory() as d:
|
||||
state = os.path.join(d, "s.json")
|
||||
save_state(state, 3)
|
||||
bridge, link, _ = self.make(rows, state)
|
||||
asyncio.run(bridge.fetch())
|
||||
self.assertEqual([m.id for m, _ in self.drain(link)], [4, 5, 6, 7])
|
||||
self.assertEqual(load_state(state), 7)
|
||||
# A second read (e.g. the WebSocket woke it up again) repeats nothing
|
||||
asyncio.run(bridge.fetch())
|
||||
self.assertEqual(self.drain(link), [])
|
||||
# New messages after a restart (a new Bridge reading the saved state)
|
||||
rows.append(row(8))
|
||||
bridge2, link2, _ = self.make(rows, state)
|
||||
asyncio.run(bridge2.fetch())
|
||||
self.assertEqual([m.id for m, _ in self.drain(link2)], [8])
|
||||
|
||||
def test_filtered_rows_still_move_the_cursor(self):
|
||||
rows = [row(1, channel="whisper"), row(2, blocked=True), row(3, channel="web", sender="[Discord] me")]
|
||||
with tempfile.TemporaryDirectory() as d:
|
||||
# Two channels: the server sends everything and the bridge drops what it mustn't relay
|
||||
bridge, link, _ = self.make(rows, os.path.join(d, "s.json"), game=GameFilter(channels=frozenset({"zone", "web"})))
|
||||
asyncio.run(bridge.fetch())
|
||||
self.assertEqual(self.drain(link), [])
|
||||
self.assertEqual(bridge.cursor.last_id, 3)
|
||||
|
||||
def test_pages_through_a_backlog(self):
|
||||
rows = [row(i) for i in range(1, 1201)]
|
||||
with tempfile.TemporaryDirectory() as d:
|
||||
state = os.path.join(d, "s.json")
|
||||
save_state(state, 1)
|
||||
bridge, link, _ = self.make(rows, state)
|
||||
link.queue = asyncio.Queue() # unbounded, to count them all
|
||||
asyncio.run(bridge.fetch())
|
||||
self.assertEqual(len(self.drain(link)), 1199)
|
||||
|
||||
def test_rate_limit_to_remote(self):
|
||||
rows = [row(i) for i in range(1, 30)]
|
||||
with tempfile.TemporaryDirectory() as d:
|
||||
state = os.path.join(d, "s.json")
|
||||
save_state(state, 1)
|
||||
bridge, link, _ = self.make(rows, state, per_min=6)
|
||||
asyncio.run(bridge.fetch())
|
||||
self.assertEqual(len(self.drain(link)), 5) # the burst
|
||||
self.assertGreater(link.out_limit.dropped, 0)
|
||||
|
||||
def test_deliver_posts_with_label(self):
|
||||
with tempfile.TemporaryDirectory() as d:
|
||||
bridge, link, dash = self.make([], os.path.join(d, "s.json"), game=GameFilter(send_zone=1200))
|
||||
asyncio.run(bridge.deliverer(link)(RemoteMessage("Alice", "hey")))
|
||||
self.assertEqual(dash.sent, [{"message": "hey", "name": "Alice", "label": "Discord", "zone": 1200}])
|
||||
|
||||
@staticmethod
|
||||
def drain(link):
|
||||
out = []
|
||||
while not link.queue.empty():
|
||||
out.append(link.queue.get_nowait())
|
||||
return out
|
||||
|
||||
|
||||
class Limits(unittest.TestCase):
|
||||
def test_token_bucket(self):
|
||||
r = RateLimiter(per_minute=60, burst=2)
|
||||
self.assertTrue(r.allow(0.0))
|
||||
self.assertTrue(r.allow(0.0))
|
||||
self.assertFalse(r.allow(0.0))
|
||||
self.assertTrue(r.allow(1.0)) # one a second comes back
|
||||
self.assertEqual(r.take_dropped(), 1)
|
||||
self.assertEqual(r.dropped, 0)
|
||||
|
||||
def test_zero_is_unlimited(self):
|
||||
r = RateLimiter(per_minute=0)
|
||||
self.assertTrue(all(r.allow(0.0) for _ in range(100)))
|
||||
|
||||
|
||||
class ConfigParsing(unittest.TestCase):
|
||||
def test_env_expansion(self):
|
||||
env = {"TOK": "abc"}
|
||||
self.assertEqual(expand_env({"a": "$TOK", "b": ["${TOK}", "x$TOK"]}, env), {"a": "abc", "b": ["abc", "x$TOK"]})
|
||||
|
||||
def test_parse(self):
|
||||
cfg = parse_config({"dashboard": {"token": "$T"}, "links": [{"adapter": "console", "console": {"name": "N"}}]},
|
||||
environ={"T": "tok"}, default_labels={"console": "Console"})
|
||||
self.assertEqual(cfg.token, "tok")
|
||||
self.assertEqual(cfg.links[0].label, "Console")
|
||||
self.assertEqual(cfg.links[0].game.channels, frozenset({"zone"}))
|
||||
self.assertEqual(cfg.links[0].settings, {"name": "N"})
|
||||
|
||||
def test_env_overrides_and_errors(self):
|
||||
cfg = parse_config({"links": [{"adapter": "console"}]}, environ={"DLU_API_TOKEN": "x", "DLU_DASHBOARD_URL": "https://d"})
|
||||
self.assertEqual((cfg.url, cfg.token), ("https://d", "x"))
|
||||
with self.assertRaises(ValueError):
|
||||
parse_config({"links": [{"adapter": "console"}]}, environ={})
|
||||
with self.assertRaises(ValueError):
|
||||
parse_config({"dashboard": {"token": "x"}, "links": [{"adapter": "console", "label": "[bad]"}]}, environ={})
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
80
tools/chat-bridge/tests/test_discord.py
Normal file
80
tools/chat-bridge/tests/test_discord.py
Normal file
@@ -0,0 +1,80 @@
|
||||
import os
|
||||
import sys
|
||||
import unittest
|
||||
|
||||
sys.path.insert(0, os.path.join(os.path.dirname(__file__), ".."))
|
||||
|
||||
from chat_bridge.adapters.discord import (DiscordAdapter, display_name, escape_markdown, from_discord, # noqa: E402
|
||||
should_forward, to_discord, webhook_username)
|
||||
from chat_bridge.core import GameMessage # noqa: E402
|
||||
|
||||
|
||||
def game(**kw):
|
||||
base = dict(id=1, channel="zone", sender_name="Bob", message="hi", zone_id=1200, zone_name="Nimbus Station")
|
||||
base.update(kw)
|
||||
return GameMessage(**base)
|
||||
|
||||
|
||||
class ToDiscord(unittest.TestCase):
|
||||
def test_bot_message(self):
|
||||
body = to_discord(game(message="*bold* @everyone"), via_webhook=False)
|
||||
self.assertEqual(body["content"], "**Bob** (Nimbus Station): \\*bold\\* @everyone")
|
||||
self.assertEqual(body["allowed_mentions"], {"parse": []}) # nobody is pinged
|
||||
|
||||
def test_webhook_message(self):
|
||||
body = to_discord(game(sender_name="DiscordFan", message="`x`"), via_webhook=True)
|
||||
self.assertEqual(body["username"], "D1scordFan (Nimbus Station)")
|
||||
self.assertEqual(body["content"], "\\`x\\`")
|
||||
self.assertEqual(body["allowed_mentions"], {"parse": []})
|
||||
|
||||
def test_no_zone(self):
|
||||
self.assertEqual(to_discord(game(zone_id=0, zone_name=""), False)["content"], "**Bob**: hi")
|
||||
|
||||
def test_skipped_note(self):
|
||||
self.assertTrue(to_discord(game(), False, skipped=4)["content"].startswith("-# 4 message(s) skipped"))
|
||||
|
||||
def test_escape(self):
|
||||
self.assertEqual(escape_markdown("a_b|c>d[e](f)"), "a\\_b\\|c\\>d\\[e\\](f)")
|
||||
|
||||
def test_webhook_username(self):
|
||||
self.assertEqual(webhook_username("clyde"), "Clyd3")
|
||||
self.assertEqual(webhook_username(" "), "Player")
|
||||
self.assertEqual(len(webhook_username("x" * 200)), 80)
|
||||
|
||||
|
||||
class FromDiscord(unittest.TestCase):
|
||||
def event(self, **kw):
|
||||
base = {"type": 0, "channel_id": "55", "content": "hello", "author": {"id": "1", "username": "alice", "global_name": "Alice"},
|
||||
"member": {"nick": None}, "mentions": [], "attachments": []}
|
||||
base.update(kw)
|
||||
return base
|
||||
|
||||
def test_display_name(self):
|
||||
self.assertEqual(display_name(self.event()), "Alice")
|
||||
self.assertEqual(display_name(self.event(member={"nick": "Ally"})), "Ally")
|
||||
self.assertEqual(display_name(self.event(author={"id": "1", "username": "alice"})), "alice")
|
||||
|
||||
def test_text(self):
|
||||
e = self.event(content="hi <@2> and <@!3> in <#9> <:lego:123> <@&4> <t:1700000000:R>",
|
||||
mentions=[{"id": "2", "username": "bob"}, {"id": "3", "username": "c", "member": {"nick": "Cee"}}],
|
||||
attachments=[{"id": "x"}])
|
||||
self.assertEqual(from_discord(e), "hi @bob and @Cee in #channel :lego: @role (a time) [attachment]")
|
||||
|
||||
def test_forward_rules(self):
|
||||
self.assertTrue(should_forward(self.event(), "55", "999"))
|
||||
self.assertFalse(should_forward(self.event(channel_id="56"), "55", "999")) # another channel
|
||||
self.assertFalse(should_forward(self.event(webhook_id="7"), "55", "999")) # our own webhook posts
|
||||
self.assertFalse(should_forward(self.event(author={"id": "8", "bot": True}), "55", "999"))
|
||||
self.assertFalse(should_forward(self.event(author={"id": "999"}), "55", "999")) # the bot itself
|
||||
self.assertFalse(should_forward(self.event(type=7), "55", "999")) # a join notice
|
||||
self.assertTrue(should_forward(self.event(type=19), "55", "999")) # a reply
|
||||
|
||||
def test_config_validation(self):
|
||||
with self.assertRaises(ValueError):
|
||||
DiscordAdapter({})
|
||||
DiscordAdapter({"webhook_url": "https://example.invalid/hook"})
|
||||
DiscordAdapter({"bot_token": "t", "channel_id": "1"})
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
Reference in New Issue
Block a user