Compare commits
17 Commits
fa97e4b098
...
v0.6.1
| Author | SHA1 | Date | |
|---|---|---|---|
| b22e7b4903 | |||
| 2867582c01 | |||
| 16b849ebd8 | |||
| 35bd68516a | |||
| 86257a70e3 | |||
| 43b5d8e3a6 | |||
| 0e9b621895 | |||
| f3f1e24c0b | |||
| 2e4228b22a | |||
| 0b79d78715 | |||
| a90cd66d3e | |||
| cb520814d8 | |||
| 4287a0d20c | |||
| e10c920a56 | |||
| 156afd6b61 | |||
| 2d47483d55 | |||
| fc4f664a5c |
1
.gitignore
vendored
1
.gitignore
vendored
@@ -6,3 +6,4 @@ dist/
|
|||||||
*.egg-info/
|
*.egg-info/
|
||||||
*.swp
|
*.swp
|
||||||
*.swo
|
*.swo
|
||||||
|
*.tmp
|
||||||
68
README.md
68
README.md
@@ -9,10 +9,12 @@ there, but it aims to be convenient and usable for relatively serious projects.
|
|||||||
The library supports the following features:
|
The library supports the following features:
|
||||||
- **Completely `asyncio` based**
|
- **Completely `asyncio` based**
|
||||||
- **Filter-based callback system**
|
- **Filter-based callback system**
|
||||||
|
- **Downloading and transparently decrypting files**
|
||||||
|
- **Sending files**
|
||||||
- **Sending images**
|
- **Sending images**
|
||||||
- **Sending videos with automatic thumbnail generation (requires `ffmpeg`)**
|
- **Sending videos with automatic thumbnail generation (requires `ffmpeg`)**
|
||||||
|
|
||||||
## 🚀 Usage
|
## 📦 Installation
|
||||||
|
|
||||||
Use `apt` to install required system packages and `pip` to install the package.
|
Use `apt` to install required system packages and `pip` to install the package.
|
||||||
You may need to use `root` privileges to use `apt`. It's highly recommended you
|
You may need to use `root` privileges to use `apt`. It's highly recommended you
|
||||||
@@ -21,18 +23,78 @@ install the latest version of the library:
|
|||||||
|
|
||||||
```bash
|
```bash
|
||||||
apt install libmagic1-dev libolm-dev
|
apt install libmagic1-dev libolm-dev
|
||||||
python -m pip install git+https://git.tyukalov.su/nikita/mab@v0.4.0
|
python -m pip install git+https://git.tyukalov.su/nikita/mab@v0.6.1
|
||||||
```
|
```
|
||||||
|
|
||||||
`libmagic1-dev` is needed for automatic file MIME type detection, `libolm-dev`
|
`libmagic1-dev` is needed for automatic file MIME type detection, `libolm-dev`
|
||||||
is needed for E2EE to work.
|
is needed for E2EE to work.
|
||||||
|
|
||||||
Please inspect [`examples/shell_bot.py`](examples/shell_bot.py),
|
Please inspect [`examples/image_bot.py`](examples/image_bot.py),
|
||||||
[`examples/echo_bot.py`](examples/echo_bot.py) or open [`examples/`](examples/)
|
[`examples/echo_bot.py`](examples/echo_bot.py) or open [`examples/`](examples/)
|
||||||
directory to find usage examples. Examples require that you set
|
directory to find usage examples. Examples require that you set
|
||||||
`MATRIX_HOMESERVER` and `MATRIX_USERNAME` environment variables. Examples create
|
`MATRIX_HOMESERVER` and `MATRIX_USERNAME` environment variables. Examples create
|
||||||
`session_storage` directory in working directory.
|
`session_storage` directory in working directory.
|
||||||
|
|
||||||
|
## 🚀 Usage
|
||||||
|
|
||||||
|
If you use `mab`, your application will *most likely* be using **callbacks** to
|
||||||
|
react to user actions. `mab` uses filter-based callback system to avoid exposing
|
||||||
|
raw `nio-matrix` event objects.
|
||||||
|
|
||||||
|
This is the workflow you will most likely follow:
|
||||||
|
1. **Define the callback as `async` function that take 1 argument of type
|
||||||
|
`EventContext`.** For example, this callback would print the caption of the
|
||||||
|
message:
|
||||||
|
```python
|
||||||
|
from mab import *
|
||||||
|
|
||||||
|
async def on_media_with_body(ctx: EventContext):
|
||||||
|
"""To be called when a message with image/video and caption is received."""
|
||||||
|
print(ctx[CTX_BODY])
|
||||||
|
```
|
||||||
|
2. **Define the conditions your callback must be called on.** For example, you
|
||||||
|
may want your callback to be called when `the sender is not the bot` and
|
||||||
|
`the message contains textual body` and (`the message is an image` or
|
||||||
|
`the message is a video`).
|
||||||
|
3. **Define the conditions as `filters`.** Most of them are pretty
|
||||||
|
straightforward. For example, if you want to use the conditions from above:
|
||||||
|
```python
|
||||||
|
from mab import *
|
||||||
|
|
||||||
|
filters = (
|
||||||
|
~SenderIsBotFilter()
|
||||||
|
& MessageTypeFilter([MessageType.IMAGE, MessageType.VIDEO])
|
||||||
|
& BodyExistsFilter()
|
||||||
|
)
|
||||||
|
```
|
||||||
|
4. **Add the callback to your `MatrixBot` instance.** For example, if you would
|
||||||
|
have used everything from above, then your code would look something like
|
||||||
|
this:
|
||||||
|
```python
|
||||||
|
from mab import *
|
||||||
|
|
||||||
|
# let's assume you create your MatrixBot as `bot` variable here
|
||||||
|
|
||||||
|
async def on_media_with_body(ctx: EventContext):
|
||||||
|
"""To be called when a message with image/video and caption is received."""
|
||||||
|
print(ctx[CTX_BODY])
|
||||||
|
|
||||||
|
filters = (
|
||||||
|
~SenderIsBotFilter()
|
||||||
|
& MessageTypeFilter([MessageType.IMAGE, MessageType.VIDEO])
|
||||||
|
& BodyExistsFilter()
|
||||||
|
)
|
||||||
|
bot.add_callback(filters, on_media_with_body)
|
||||||
|
|
||||||
|
...
|
||||||
|
```
|
||||||
|
|
||||||
|
Filters support bitwise operators to implement complex matching logic. Some
|
||||||
|
filters set context variables which can be accessed by
|
||||||
|
`context[CTX_KEY_NAME]`-like syntax. Possible variables are defined in
|
||||||
|
[this file](src/mab/context.py). Filters are implemented in
|
||||||
|
[files of this directory](src/mab/filters/).
|
||||||
|
|
||||||
## 🏷️ Versioning
|
## 🏷️ Versioning
|
||||||
|
|
||||||
Releases are tagged in this repository using the `vX.Y.Z` format. If the commit
|
Releases are tagged in this repository using the `vX.Y.Z` format. If the commit
|
||||||
|
|||||||
130
examples/command_bot.py
Normal file
130
examples/command_bot.py
Normal file
@@ -0,0 +1,130 @@
|
|||||||
|
"""
|
||||||
|
This example implements Matrix bot that can execute some commands.
|
||||||
|
|
||||||
|
It uses environment variables to specify authorization data. Use Ctrl+C to stop
|
||||||
|
the bot.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
import time
|
||||||
|
import html
|
||||||
|
import os
|
||||||
|
import traceback
|
||||||
|
import logging
|
||||||
|
|
||||||
|
from mab import *
|
||||||
|
|
||||||
|
from _environment import check_environment
|
||||||
|
|
||||||
|
async def on_help_command(ctx: EventContext) -> None:
|
||||||
|
"""!help"""
|
||||||
|
HELP_MESSAGE = (
|
||||||
|
"<strong>Here is the list of the commands:</strong><br>"
|
||||||
|
"<ul>"
|
||||||
|
"<li><code>!help</code> - this help message</li>"
|
||||||
|
"<li><code>!time</code> - get UNIX timestamp</li>"
|
||||||
|
"<li><code>!raise</code> - raise <code>RuntimeError()</code></li>"
|
||||||
|
"<li><code>!assert</code> - perform <code>assert</code> that will fail</li>"
|
||||||
|
"<li><code>!mul A B [C] [D]...</code> - multiply A, B... and so on</li>"
|
||||||
|
"<li><code>!args arg1 [arg2] ... [arg5]</code> - command that takes 1..5 arguments</li>"
|
||||||
|
"</ul>"
|
||||||
|
)
|
||||||
|
await ctx.bot.send.text(ctx.room, HELP_MESSAGE)
|
||||||
|
|
||||||
|
async def on_time_command(ctx: EventContext) -> None:
|
||||||
|
"""!time"""
|
||||||
|
await ctx.bot.send.text(
|
||||||
|
ctx.room,
|
||||||
|
f"Current UNIX timestamp is <strong>{int(time.time())}</strong>"
|
||||||
|
)
|
||||||
|
|
||||||
|
async def on_raise_command(ctx: EventContext) -> None:
|
||||||
|
"""!raise"""
|
||||||
|
await ctx.bot.send.text(
|
||||||
|
ctx.room,
|
||||||
|
"<strong>Executing <code>raise RuntimeError()</code>...</strong>"
|
||||||
|
)
|
||||||
|
raise RuntimeError()
|
||||||
|
|
||||||
|
async def on_assert_command(ctx: EventContext) -> None:
|
||||||
|
"""!assert"""
|
||||||
|
await ctx.bot.send.text(
|
||||||
|
ctx.room,
|
||||||
|
"<strong>Executing <code>assert False</code>...</strong>"
|
||||||
|
)
|
||||||
|
assert False
|
||||||
|
|
||||||
|
async def on_mul_command(ctx: EventContext) -> None:
|
||||||
|
"""!mul"""
|
||||||
|
try:
|
||||||
|
numbers = [float(v) for v in ctx[CTX_CMD_ARGS]]
|
||||||
|
v = numbers[0]
|
||||||
|
for n in numbers[1:]:
|
||||||
|
v *= n
|
||||||
|
response = " * ".join(html.escape("%.2f" % n) for n in numbers)
|
||||||
|
response += f" = <strong>{html.escape(str(v))}<strong>"
|
||||||
|
await ctx.bot.send.text(ctx.room, response)
|
||||||
|
except Exception as e:
|
||||||
|
await ctx.bot.send.text(ctx.room, f"Could not process the command: {e}")
|
||||||
|
|
||||||
|
async def on_args_command(ctx: EventContext) -> None:
|
||||||
|
"""!args"""
|
||||||
|
try:
|
||||||
|
response = (
|
||||||
|
f"Prefix: <code>{ctx[CTX_CMD_PREFIX]}</code><br>"
|
||||||
|
f"Verb: <code>{ctx[CTX_CMD_VERB]}</code><br>"
|
||||||
|
f"Arguments: <code>{len(ctx[CTX_CMD_ARGS])}</code><br>"
|
||||||
|
f"Arguments are:<br><ol>"
|
||||||
|
)
|
||||||
|
for arg in ctx[CTX_CMD_ARGS]:
|
||||||
|
response += f"<li><code>{html.escape(arg)}</code></li>"
|
||||||
|
response += "</ol>"
|
||||||
|
await ctx.bot.send.text(
|
||||||
|
ctx.room,
|
||||||
|
response
|
||||||
|
)
|
||||||
|
except:
|
||||||
|
await ctx.bot.send.text(
|
||||||
|
ctx.room,
|
||||||
|
f"Could not process the command: {traceback.format_exc()}"
|
||||||
|
)
|
||||||
|
|
||||||
|
async def invalid_usage(ctx: EventContext) -> None:
|
||||||
|
"""This callback is called when the bot used incorrectly."""
|
||||||
|
await ctx.bot.send.text(ctx.room, "Use <code>!help</code>")
|
||||||
|
|
||||||
|
|
||||||
|
async def main() -> None:
|
||||||
|
"""Application entry point"""
|
||||||
|
logging.basicConfig(level=logging.INFO)
|
||||||
|
logging.getLogger("nio").setLevel(logging.CRITICAL+1)
|
||||||
|
check_environment()
|
||||||
|
|
||||||
|
config = MatrixBotConfig(
|
||||||
|
matrix_homeserver_url=os.environ["MATRIX_HOMESERVER"],
|
||||||
|
matrix_username_localpart=os.environ["MATRIX_USERNAME"],
|
||||||
|
storage_directory="session_storage"
|
||||||
|
)
|
||||||
|
bot = MatrixBot(config)
|
||||||
|
|
||||||
|
COMMANDS = {
|
||||||
|
on_help_command: BodyCommandFilter(["help", "?"]),
|
||||||
|
on_time_command: BodyCommandFilter("time"),
|
||||||
|
on_raise_command: BodyCommandFilter("raise"),
|
||||||
|
on_assert_command: BodyCommandFilter("assert"),
|
||||||
|
on_mul_command: BodyCommandFilter("mul", min_args=2),
|
||||||
|
on_args_command: BodyCommandFilter("args", min_args=1, max_args=5),
|
||||||
|
}
|
||||||
|
for callback, filter in COMMANDS.items():
|
||||||
|
f = ~SenderIsBotFilter() & filter
|
||||||
|
bot.add_callback(f, callback)
|
||||||
|
bot.add_callback(~SenderIsBotFilter() & NewMessageFilter(), invalid_usage)
|
||||||
|
|
||||||
|
# run until Ctrl+C
|
||||||
|
try:
|
||||||
|
await bot.run()
|
||||||
|
except asyncio.CancelledError:
|
||||||
|
pass
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
asyncio.run(main())
|
||||||
@@ -9,21 +9,13 @@ import asyncio
|
|||||||
import os
|
import os
|
||||||
import logging
|
import logging
|
||||||
|
|
||||||
from mab import (
|
from mab import *
|
||||||
MatrixBot,
|
|
||||||
MatrixBotConfig,
|
|
||||||
RoomEventData,
|
|
||||||
BodyExistsFilter,
|
|
||||||
MessageTypeFilter,
|
|
||||||
SenderIsBotFilter,
|
|
||||||
MessageType
|
|
||||||
)
|
|
||||||
|
|
||||||
from _environment import check_environment
|
from _environment import check_environment
|
||||||
|
|
||||||
async def on_text_message(data: RoomEventData) -> None:
|
async def on_text_message(ctx: EventContext) -> None:
|
||||||
"""This callback is called when a text message arrives."""
|
"""This callback is called when a text message arrives."""
|
||||||
await data.bot.send_text(data.room, data.event.body) # type: ignore
|
await ctx.bot.send.text(ctx.room, ctx[CTX_BODY])
|
||||||
|
|
||||||
async def main() -> None:
|
async def main() -> None:
|
||||||
"""Application entry point"""
|
"""Application entry point"""
|
||||||
|
|||||||
97
examples/file_bot.py
Normal file
97
examples/file_bot.py
Normal file
@@ -0,0 +1,97 @@
|
|||||||
|
"""
|
||||||
|
This example implements Matrix bot that calculates SHA256 for a file sent by
|
||||||
|
user.
|
||||||
|
|
||||||
|
It uses environment variables to specify authorization data. Use Ctrl+C to stop
|
||||||
|
the bot.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
import aiofiles
|
||||||
|
import os
|
||||||
|
import logging
|
||||||
|
import hashlib
|
||||||
|
|
||||||
|
from mab import *
|
||||||
|
|
||||||
|
from _environment import check_environment
|
||||||
|
|
||||||
|
async def on_text_message(ctx: EventContext) -> None:
|
||||||
|
"""This callback is called when a text message arrives."""
|
||||||
|
await ctx.bot.send.text(ctx.room, "Please send a file/image/video")
|
||||||
|
|
||||||
|
async def on_file_message(ctx: EventContext) -> None:
|
||||||
|
"""This callback is called when a file arrives."""
|
||||||
|
sha = hashlib.sha256()
|
||||||
|
does_temp_exist = False
|
||||||
|
if ctx[CTX_FILE_SIZE] > 1_000_000:
|
||||||
|
await ctx.bot.send.text(
|
||||||
|
ctx.room,
|
||||||
|
"The file is larger than 1 MB, downloading to filesystem"
|
||||||
|
)
|
||||||
|
await ctx.bot.download.file(ctx, path="temp.tmp")
|
||||||
|
async with aiofiles.open("temp.tmp", "rb") as f:
|
||||||
|
while True:
|
||||||
|
data = await f.read(64 * 1024)
|
||||||
|
if not data:
|
||||||
|
break
|
||||||
|
sha.update(data)
|
||||||
|
does_temp_exist = True
|
||||||
|
else:
|
||||||
|
await ctx.bot.send.text(
|
||||||
|
ctx.room,
|
||||||
|
"The file is smaller than 1 MB, downloading to RAM"
|
||||||
|
)
|
||||||
|
content = await ctx.bot.download.file(ctx, path=None)
|
||||||
|
sha.update(content)
|
||||||
|
# result
|
||||||
|
response = f"SHA256 for file `{ctx[CTX_FILE_NAME]}`"
|
||||||
|
response += f" ({ctx[CTX_FILE_SIZE]} bytes, {ctx[CTX_FILE_MIME]})"
|
||||||
|
await ctx.bot.send.file_bytes(
|
||||||
|
room=ctx.room,
|
||||||
|
data=sha.hexdigest().encode("utf-8"),
|
||||||
|
filename="hash of the file.txt",
|
||||||
|
mime_type="text/plain",
|
||||||
|
text=response
|
||||||
|
)
|
||||||
|
# resend the file to test uploading
|
||||||
|
if does_temp_exist:
|
||||||
|
await ctx.bot.send.file(
|
||||||
|
room=ctx.room,
|
||||||
|
path="temp.tmp",
|
||||||
|
filename=ctx[CTX_FILE_NAME],
|
||||||
|
text="This is the file you have sent, but it was reuploaded"
|
||||||
|
)
|
||||||
|
try:
|
||||||
|
os.unlink("temp.tmp")
|
||||||
|
except:
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
|
async def main() -> None:
|
||||||
|
"""Application entry point"""
|
||||||
|
logging.basicConfig(level=logging.INFO)
|
||||||
|
logging.getLogger("nio").setLevel(logging.CRITICAL+1)
|
||||||
|
check_environment()
|
||||||
|
|
||||||
|
config = MatrixBotConfig(
|
||||||
|
matrix_homeserver_url=os.environ["MATRIX_HOMESERVER"],
|
||||||
|
matrix_username_localpart=os.environ["MATRIX_USERNAME"],
|
||||||
|
storage_directory="session_storage"
|
||||||
|
)
|
||||||
|
bot = MatrixBot(config)
|
||||||
|
bot.add_callback(
|
||||||
|
~SenderIsBotFilter() & BodyExistsFilter() & MessageTypeFilter(MessageType.TEXT),
|
||||||
|
on_text_message)
|
||||||
|
bot.add_callback(
|
||||||
|
~SenderIsBotFilter() & MessageHasFile() & NewMessageFilter(),
|
||||||
|
on_file_message)
|
||||||
|
|
||||||
|
# run until Ctrl+C
|
||||||
|
try:
|
||||||
|
await bot.run()
|
||||||
|
except asyncio.CancelledError:
|
||||||
|
pass
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
asyncio.run(main())
|
||||||
@@ -9,31 +9,24 @@ the bot.
|
|||||||
import asyncio
|
import asyncio
|
||||||
from io import BytesIO
|
from io import BytesIO
|
||||||
import os
|
import os
|
||||||
|
import time
|
||||||
import logging
|
import logging
|
||||||
import random
|
import random
|
||||||
from PIL import Image
|
from PIL import Image, ImageFilter
|
||||||
|
|
||||||
from mab import (
|
from mab import *
|
||||||
MatrixBot,
|
|
||||||
MatrixBotConfig,
|
|
||||||
RoomEventData,
|
|
||||||
MessageTypeFilter,
|
|
||||||
BodyCommandFilter,
|
|
||||||
SenderIsBotFilter,
|
|
||||||
MessageType
|
|
||||||
)
|
|
||||||
|
|
||||||
from _environment import check_environment
|
from _environment import check_environment
|
||||||
|
|
||||||
async def on_gen_command(data: RoomEventData) -> None:
|
async def on_gen_command(data: EventContext) -> None:
|
||||||
"""This callback is called when `!gen R G B` command is received."""
|
"""This callback is called when `!gen R G B` command is received."""
|
||||||
# convert R, G and B to floats
|
# convert R, G and B to floats
|
||||||
try:
|
try:
|
||||||
r, g, b = [float(v) for v in data.event.command_args] # type: ignore
|
r, g, b = [float(v) for v in data[CTX_CMD_ARGS]]
|
||||||
except:
|
except:
|
||||||
await data.bot.send_text(data.room, "Invalid arguments")
|
await data.bot.send.text(data.room, "Invalid arguments")
|
||||||
return
|
return
|
||||||
await data.bot.send_text(data.room, "Generating the noise...")
|
await data.bot.send.text(data.room, "Generating the noise...")
|
||||||
# create the basic noise
|
# create the basic noise
|
||||||
img = Image.new("RGB", (16, 16))
|
img = Image.new("RGB", (16, 16))
|
||||||
for x in range(img.width):
|
for x in range(img.width):
|
||||||
@@ -51,13 +44,40 @@ async def on_gen_command(data: RoomEventData) -> None:
|
|||||||
buf.seek(0)
|
buf.seek(0)
|
||||||
buf = buf.read()
|
buf = buf.read()
|
||||||
# send
|
# send
|
||||||
await data.bot.send_image_bytes(data.room, buf, "noise.png")
|
await data.bot.send.image_bytes(data.room, buf, "noise.png")
|
||||||
|
|
||||||
async def on_wrong_message(data: RoomEventData) -> None:
|
async def on_image(ctx: EventContext) -> None:
|
||||||
|
"""This callback is called when an image is received."""
|
||||||
|
# download
|
||||||
|
s = time.time()
|
||||||
|
data = await ctx.bot.download.file(ctx)
|
||||||
|
took_time = time.time() - s
|
||||||
|
# convert to Image
|
||||||
|
with BytesIO(data) as buf:
|
||||||
|
img = Image.open(buf)
|
||||||
|
# apply effects
|
||||||
|
blur = ImageFilter.GaussianBlur(
|
||||||
|
radius=min(img.size[0] // 10, 5)
|
||||||
|
)
|
||||||
|
img = img.filter(blur)
|
||||||
|
# save to buffer
|
||||||
|
buf = BytesIO()
|
||||||
|
img.save(buf, format="PNG")
|
||||||
|
buf.seek(0)
|
||||||
|
buf = buf.read()
|
||||||
|
# send
|
||||||
|
await ctx.bot.send.image_bytes(
|
||||||
|
ctx.room,
|
||||||
|
buf,
|
||||||
|
"blurred.png",
|
||||||
|
text="Download and decryption took %.4f seconds" % took_time
|
||||||
|
)
|
||||||
|
|
||||||
|
async def on_wrong_message(data: EventContext) -> None:
|
||||||
"""This callback is called when a wrong message is received."""
|
"""This callback is called when a wrong message is received."""
|
||||||
await data.bot.send_text(
|
await data.bot.send.text(
|
||||||
data.room,
|
data.room,
|
||||||
"Text me something like <code>!gen 0.1 0.7 1.0</code>"
|
"Text me something like <code>!gen 0.1 0.7 1.0</code> or send an image to blur"
|
||||||
)
|
)
|
||||||
|
|
||||||
async def main() -> None:
|
async def main() -> None:
|
||||||
@@ -85,11 +105,17 @@ async def main() -> None:
|
|||||||
~SenderIsBotFilter() & command_filter,
|
~SenderIsBotFilter() & command_filter,
|
||||||
on_gen_command)
|
on_gen_command)
|
||||||
|
|
||||||
|
# callback for new image message
|
||||||
|
bot.add_callback(
|
||||||
|
~SenderIsBotFilter() & MessageTypeFilter(MessageType.IMAGE) & NewMessageFilter(),
|
||||||
|
on_image
|
||||||
|
)
|
||||||
|
|
||||||
# callback for message that
|
# callback for message that
|
||||||
# 1. are sent not by this bot
|
# 1. are sent not by this bot
|
||||||
# 2. do NOT match the command filter
|
# 2. are new messages (not edits)
|
||||||
bot.add_callback(
|
bot.add_callback(
|
||||||
~SenderIsBotFilter() & ~command_filter,
|
~SenderIsBotFilter() & BodyExistsFilter() & NewMessageFilter(),
|
||||||
on_wrong_message
|
on_wrong_message
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
|
|||||||
|
|
||||||
[project]
|
[project]
|
||||||
name = "mab"
|
name = "mab"
|
||||||
version = "0.4.0"
|
version = "0.6.1"
|
||||||
authors = [
|
authors = [
|
||||||
{ name = "Tyukalov Nikita", email = "nikita@tyukalov.su" }
|
{ name = "Tyukalov Nikita", email = "nikita@tyukalov.su" }
|
||||||
]
|
]
|
||||||
|
|||||||
@@ -1,7 +1,8 @@
|
|||||||
from . import bot
|
from . import bot
|
||||||
from . import types
|
from . import types
|
||||||
|
|
||||||
from .types import MatrixBotConfig, RoomEventData, MessageType
|
from .types import MatrixBotConfig, MessageType
|
||||||
|
from .context import *
|
||||||
|
|
||||||
from .bot import MatrixBot
|
from .bot import MatrixBot
|
||||||
|
|
||||||
@@ -16,9 +17,20 @@ __all__ = [
|
|||||||
|
|
||||||
# .types
|
# .types
|
||||||
"MatrixBotConfig",
|
"MatrixBotConfig",
|
||||||
"RoomEventData",
|
|
||||||
"MessageType",
|
"MessageType",
|
||||||
|
|
||||||
|
# .context
|
||||||
|
"EventContext",
|
||||||
|
"CTX_BODY",
|
||||||
|
"CTX_MESSAGE_TYPE",
|
||||||
|
"CTX_SENDER",
|
||||||
|
"CTX_CMD_PREFIX",
|
||||||
|
"CTX_CMD_VERB",
|
||||||
|
"CTX_CMD_ARGS",
|
||||||
|
"CTX_FILE_SIZE",
|
||||||
|
"CTX_FILE_MIME",
|
||||||
|
"CTX_FILE_NAME",
|
||||||
|
|
||||||
# .bot
|
# .bot
|
||||||
"MatrixBot",
|
"MatrixBot",
|
||||||
|
|
||||||
@@ -41,4 +53,5 @@ __all__ = [
|
|||||||
"RedactedMessageFilter",
|
"RedactedMessageFilter",
|
||||||
"SenderIsFilter",
|
"SenderIsFilter",
|
||||||
"SenderIsBotFilter",
|
"SenderIsBotFilter",
|
||||||
|
"MessageHasFile",
|
||||||
]
|
]
|
||||||
@@ -10,7 +10,8 @@ from nio import MatrixInvitedRoom, InviteMemberEvent, JoinResponse
|
|||||||
from nio.events.room_events import Event as RoomEvent
|
from nio.events.room_events import Event as RoomEvent
|
||||||
|
|
||||||
from ._storage import Storage
|
from ._storage import Storage
|
||||||
from ..types import MatrixBotConfig, RoomEventData
|
from ..types import MatrixBotConfig
|
||||||
|
from ..context import EventContext
|
||||||
from ..filters.base import BaseEventFilter
|
from ..filters.base import BaseEventFilter
|
||||||
|
|
||||||
if TYPE_CHECKING:
|
if TYPE_CHECKING:
|
||||||
@@ -31,7 +32,7 @@ class Callbacks:
|
|||||||
filter: BaseEventFilter
|
filter: BaseEventFilter
|
||||||
"""Filter to use for matching"""
|
"""Filter to use for matching"""
|
||||||
|
|
||||||
callback: Callable[[RoomEventData], Coroutine[Any, Any, None]] | None
|
callback: Callable[[EventContext], Coroutine[Any, Any, None]] | None
|
||||||
"""Callback that will be called if the filter matches"""
|
"""Callback that will be called if the filter matches"""
|
||||||
|
|
||||||
stop_matching: bool
|
stop_matching: bool
|
||||||
@@ -51,20 +52,19 @@ class Callbacks:
|
|||||||
for callback_info in self._filters:
|
for callback_info in self._filters:
|
||||||
if not isinstance(callback_info, self._FilterBasedCallback):
|
if not isinstance(callback_info, self._FilterBasedCallback):
|
||||||
continue
|
continue
|
||||||
|
event_data = EventContext(
|
||||||
|
room=room,
|
||||||
|
event=event,
|
||||||
|
bot=self._matrix_bot
|
||||||
|
)
|
||||||
try:
|
try:
|
||||||
if not await callback_info.filter(room, event, self._client):
|
if not await callback_info.filter(event_data):
|
||||||
continue
|
continue
|
||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
raise
|
raise
|
||||||
except:
|
except:
|
||||||
self._logger.error(traceback.format_exc())
|
self._logger.error(traceback.format_exc())
|
||||||
continue
|
continue
|
||||||
event_data = RoomEventData(
|
|
||||||
room=room,
|
|
||||||
event=event,
|
|
||||||
filter=callback_info.filter,
|
|
||||||
bot=self._matrix_bot
|
|
||||||
)
|
|
||||||
try:
|
try:
|
||||||
# dump argument types
|
# dump argument types
|
||||||
if callback_info.callback is None:
|
if callback_info.callback is None:
|
||||||
@@ -144,7 +144,7 @@ class Callbacks:
|
|||||||
def add_room_event_callback(
|
def add_room_event_callback(
|
||||||
self,
|
self,
|
||||||
filter: BaseEventFilter,
|
filter: BaseEventFilter,
|
||||||
callback: Callable[[RoomEventData], Coroutine[Any, Any, None]] | None,
|
callback: Callable[[EventContext], Coroutine[Any, Any, None]] | None,
|
||||||
*,
|
*,
|
||||||
stop_matching: bool = True) -> None:
|
stop_matching: bool = True) -> None:
|
||||||
"""
|
"""
|
||||||
|
|||||||
250
src/mab/bot/_client_downloader.py
Normal file
250
src/mab/bot/_client_downloader.py
Normal file
@@ -0,0 +1,250 @@
|
|||||||
|
import asyncio
|
||||||
|
import aiofiles
|
||||||
|
import logging
|
||||||
|
import os
|
||||||
|
from typing import overload
|
||||||
|
import threading
|
||||||
|
from uuid import uuid4
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
import hashlib
|
||||||
|
import hmac
|
||||||
|
import unpaddedbase64
|
||||||
|
from Crypto.Cipher import AES
|
||||||
|
from Crypto.Util import Counter
|
||||||
|
|
||||||
|
from nio import (
|
||||||
|
AsyncClient,
|
||||||
|
Event,
|
||||||
|
DiskDownloadResponse,
|
||||||
|
MemoryDownloadResponse
|
||||||
|
)
|
||||||
|
|
||||||
|
from ..types import MatrixBotConfig
|
||||||
|
from ..context import EventContext
|
||||||
|
|
||||||
|
class ClientDownloader:
|
||||||
|
"""This class downloads files"""
|
||||||
|
|
||||||
|
def __init__(self):
|
||||||
|
self._logger = logging.getLogger("ClientSender")
|
||||||
|
self._config: MatrixBotConfig | None = None
|
||||||
|
self._client: AsyncClient | None = None
|
||||||
|
|
||||||
|
async def setup(self,
|
||||||
|
config: MatrixBotConfig,
|
||||||
|
client: AsyncClient) -> None:
|
||||||
|
"""
|
||||||
|
Setup the downloader.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
- config - config to use
|
||||||
|
- client - client to use
|
||||||
|
"""
|
||||||
|
self._config = config
|
||||||
|
self._client = client
|
||||||
|
|
||||||
|
async def _decrypt_file(self,
|
||||||
|
src: Path,
|
||||||
|
dst: Path,
|
||||||
|
key: str,
|
||||||
|
iv: str,
|
||||||
|
sha256: str) -> None:
|
||||||
|
"""
|
||||||
|
Decrypts `src` and saves it as `dst` asynchronously.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
- src - path to the source (encrypted) file
|
||||||
|
- dst - path where the resulting (decrypted) file will be saved
|
||||||
|
- key - key, from event["content"]["file"]["key"]["k"]
|
||||||
|
- iv - initialization vector, from event["content"]["file"]["iv"]
|
||||||
|
- sha256 - SHA-256 digest, from
|
||||||
|
event["content"]["file"]["hashes"]["sha256"]
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
- Returns nothing on success
|
||||||
|
- Raises an exception on error
|
||||||
|
"""
|
||||||
|
async with aiofiles.open(src, "rb") as reader:
|
||||||
|
# check file hash
|
||||||
|
expected = unpaddedbase64.decode_base64(sha256)
|
||||||
|
digest = hashlib.sha256()
|
||||||
|
while chunk := await reader.read(16 * 1024):
|
||||||
|
digest.update(chunk)
|
||||||
|
if not hmac.compare_digest(digest.digest(), expected):
|
||||||
|
raise RuntimeError("SHA256 mismatch")
|
||||||
|
await reader.seek(0)
|
||||||
|
# decrypt
|
||||||
|
decoded_key = unpaddedbase64.decode_base64(key)
|
||||||
|
decoded_iv = unpaddedbase64.decode_base64(iv)
|
||||||
|
cipher = AES.new(
|
||||||
|
decoded_key,
|
||||||
|
AES.MODE_CTR,
|
||||||
|
counter=Counter.new(
|
||||||
|
nbits=64,
|
||||||
|
prefix=decoded_iv[:8],
|
||||||
|
initial_value=int.from_bytes(decoded_iv[8:], "big")
|
||||||
|
)
|
||||||
|
)
|
||||||
|
async with aiofiles.open(dst, "wb") as writer:
|
||||||
|
while chunk := await reader.read(16 * 1024):
|
||||||
|
await writer.write(cipher.decrypt(chunk))
|
||||||
|
|
||||||
|
async def _decrypt_bytes(self,
|
||||||
|
data: bytes,
|
||||||
|
key: str,
|
||||||
|
iv: str,
|
||||||
|
sha256: str) -> bytes:
|
||||||
|
"""
|
||||||
|
Decrypts `data`. It starts subthread that performs decryption.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
- data - data that needs to be decrypted
|
||||||
|
- key - key, from event["content"]["file"]["key"]["k"]
|
||||||
|
- iv - initialization vector, from event["content"]["file"]["iv"]
|
||||||
|
- sha265 - SHA-256 digest, from
|
||||||
|
event["content"]["file"]["hashes"]["sha256"]
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
- Returns decrypted data on success
|
||||||
|
- Raises an exception on error
|
||||||
|
"""
|
||||||
|
cancel = threading.Event()
|
||||||
|
# subfunction that will run in another thread
|
||||||
|
def decrypt() -> bytes:
|
||||||
|
chunk_size = 64 * 1024
|
||||||
|
view = memoryview(data)
|
||||||
|
expected = unpaddedbase64.decode_base64(sha256)
|
||||||
|
digest = hashlib.sha256()
|
||||||
|
for offset in range(0, len(view), chunk_size):
|
||||||
|
if cancel.is_set():
|
||||||
|
return b""
|
||||||
|
digest.update(view[offset:offset + chunk_size])
|
||||||
|
if not hmac.compare_digest(digest.digest(), expected):
|
||||||
|
raise RuntimeError("SHA-256 mismatch")
|
||||||
|
# prepare AES-CTR
|
||||||
|
decoded_key = unpaddedbase64.decode_base64(key)
|
||||||
|
decoded_iv = unpaddedbase64.decode_base64(iv)
|
||||||
|
cipher = AES.new(
|
||||||
|
decoded_key,
|
||||||
|
AES.MODE_CTR,
|
||||||
|
counter=Counter.new(
|
||||||
|
nbits=64,
|
||||||
|
prefix=decoded_iv[:8],
|
||||||
|
initial_value=int.from_bytes(decoded_iv[8:], "big"),
|
||||||
|
),
|
||||||
|
)
|
||||||
|
# decrypt
|
||||||
|
result = bytearray()
|
||||||
|
for offset in range(0, len(view), chunk_size):
|
||||||
|
if cancel.is_set():
|
||||||
|
return b""
|
||||||
|
result.extend(
|
||||||
|
cipher.decrypt(view[offset:offset + chunk_size])
|
||||||
|
)
|
||||||
|
return bytes(result)
|
||||||
|
# decrypt in thread
|
||||||
|
try:
|
||||||
|
return await asyncio.to_thread(decrypt)
|
||||||
|
except asyncio.CancelledError:
|
||||||
|
cancel.set()
|
||||||
|
raise
|
||||||
|
|
||||||
|
@overload
|
||||||
|
async def file(self,
|
||||||
|
source: EventContext | Event,
|
||||||
|
*,
|
||||||
|
path: str | Path) -> Path: ...
|
||||||
|
|
||||||
|
@overload
|
||||||
|
async def file(self,
|
||||||
|
source: EventContext | Event,
|
||||||
|
*,
|
||||||
|
path: None = None) -> bytes: ...
|
||||||
|
|
||||||
|
async def file(self,
|
||||||
|
source: EventContext | Event,
|
||||||
|
*,
|
||||||
|
path: str | Path | None = None) -> Path | bytes:
|
||||||
|
"""
|
||||||
|
Download the file from the event. Automatically deciphers encrypted
|
||||||
|
media.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
- source - event that has a file in its `content`
|
||||||
|
- path - where to save the file to. Use `None` to store the file in
|
||||||
|
memory. Use path to specify file download path.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
- Returns `Path` to the file on success (if `path` isn't `None`)
|
||||||
|
- Returns `bytes` of the file on success (if `path` is `None`)
|
||||||
|
- Raises an exception on error
|
||||||
|
"""
|
||||||
|
if self._client is None:
|
||||||
|
raise RuntimeError("ClientDownloader is not set up")
|
||||||
|
# prepare the args
|
||||||
|
if isinstance(source, EventContext):
|
||||||
|
source = source.event
|
||||||
|
content: dict = source.source["content"]
|
||||||
|
f = content.get("file")
|
||||||
|
# encypted
|
||||||
|
if f:
|
||||||
|
if not isinstance(f["url"], str):
|
||||||
|
raise RuntimeError("`source` does not contain file URL")
|
||||||
|
mxc = f["url"]
|
||||||
|
if f["key"]["alg"].upper() != "A256CTR":
|
||||||
|
raise NotImplementedError(
|
||||||
|
f"Unsupported encryption algorithm `{f['key']['alg']}`"
|
||||||
|
)
|
||||||
|
key: str | None = f["key"]["k"]
|
||||||
|
iv: str | None = f["iv"]
|
||||||
|
sha256: str | None = f["hashes"]["sha256"]
|
||||||
|
# not encrypted
|
||||||
|
else:
|
||||||
|
mxc = content["url"]
|
||||||
|
key = None
|
||||||
|
iv = None
|
||||||
|
sha256 = None
|
||||||
|
# prepare path
|
||||||
|
if isinstance(path, str):
|
||||||
|
path = Path(path)
|
||||||
|
filename: str | None = \
|
||||||
|
os.path.basename(path) if isinstance(path, Path) else None
|
||||||
|
# download
|
||||||
|
result = await self._client.download(
|
||||||
|
mxc,
|
||||||
|
filename=filename,
|
||||||
|
save_to=path
|
||||||
|
)
|
||||||
|
# check the response
|
||||||
|
if isinstance(result, DiskDownloadResponse):
|
||||||
|
result = Path(result.body)
|
||||||
|
elif isinstance(result, MemoryDownloadResponse):
|
||||||
|
result = result.body
|
||||||
|
else:
|
||||||
|
raise RuntimeError(result)
|
||||||
|
# decrypt if needed
|
||||||
|
temp_path: Path | None = None
|
||||||
|
try:
|
||||||
|
if key is not None and iv is not None and sha256 is not None:
|
||||||
|
# on disk
|
||||||
|
if isinstance(result, Path):
|
||||||
|
temp_path = result.with_name(
|
||||||
|
f".{result.name}.{uuid4().hex}.tmp"
|
||||||
|
)
|
||||||
|
await self._decrypt_file(
|
||||||
|
result, temp_path, key, iv, sha256
|
||||||
|
)
|
||||||
|
temp_path.replace(result)
|
||||||
|
temp_path = None
|
||||||
|
# in memory
|
||||||
|
else:
|
||||||
|
result = await self._decrypt_bytes(
|
||||||
|
result, key, iv, sha256
|
||||||
|
)
|
||||||
|
finally:
|
||||||
|
if isinstance(temp_path, Path) and isinstance(result, Path):
|
||||||
|
result.unlink(missing_ok=True)
|
||||||
|
temp_path.unlink(missing_ok=True)
|
||||||
|
# return the result
|
||||||
|
return result
|
||||||
147
src/mab/bot/_client_room_operations.py
Normal file
147
src/mab/bot/_client_room_operations.py
Normal file
@@ -0,0 +1,147 @@
|
|||||||
|
import logging
|
||||||
|
from nio import (
|
||||||
|
AsyncClient,
|
||||||
|
Event,
|
||||||
|
MatrixRoom,
|
||||||
|
MatrixUser,
|
||||||
|
RoomMember,
|
||||||
|
|
||||||
|
RoomGetEventResponse,
|
||||||
|
RoomLeaveResponse,
|
||||||
|
JoinResponse,
|
||||||
|
RoomInviteResponse,
|
||||||
|
RoomKickResponse,
|
||||||
|
RoomBanResponse,
|
||||||
|
RoomUnbanResponse,
|
||||||
|
JoinedMembersResponse
|
||||||
|
)
|
||||||
|
from ..utils import MatrixBotConfig
|
||||||
|
|
||||||
|
class ClientRoomOperations:
|
||||||
|
"""This class works with rooms."""
|
||||||
|
|
||||||
|
@staticmethod
|
||||||
|
def _room_id(room: MatrixRoom | str) -> str:
|
||||||
|
"""Convert room to its ID if needed"""
|
||||||
|
return room.room_id if isinstance(room, MatrixRoom) else room
|
||||||
|
|
||||||
|
@staticmethod
|
||||||
|
def _user_id(user: MatrixUser | str) -> str:
|
||||||
|
"""Convert user to its ID if needed"""
|
||||||
|
return user.user_id if isinstance(user, MatrixUser) else user
|
||||||
|
|
||||||
|
def __init__(self):
|
||||||
|
self._logger = logging.getLogger("ClientRoomOperations")
|
||||||
|
self._config: MatrixBotConfig | None = None
|
||||||
|
self._client: AsyncClient | None = None
|
||||||
|
|
||||||
|
async def setup(self,
|
||||||
|
config: MatrixBotConfig,
|
||||||
|
client: AsyncClient) -> None:
|
||||||
|
"""Setup room operations.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
- config - config to use
|
||||||
|
- client - client to use
|
||||||
|
"""
|
||||||
|
self._config = config
|
||||||
|
self._client = client
|
||||||
|
|
||||||
|
async def event(self, room: MatrixRoom | str, event_id: str) -> Event:
|
||||||
|
"""Get raw event information from matrix server (no cache used yet).
|
||||||
|
|
||||||
|
Args:
|
||||||
|
- room - room where the event happened
|
||||||
|
- event_id - ID of the event
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
- `nio.Event` on success
|
||||||
|
- raises an exception on matrix error
|
||||||
|
"""
|
||||||
|
if self._config is None or self._client is None:
|
||||||
|
raise RuntimeError("ClientRoom is not set up")
|
||||||
|
res = await self._client.room_get_event(
|
||||||
|
self._room_id(room),
|
||||||
|
event_id
|
||||||
|
)
|
||||||
|
if isinstance(res, RoomGetEventResponse):
|
||||||
|
return res.event
|
||||||
|
else:
|
||||||
|
raise RuntimeError(res)
|
||||||
|
|
||||||
|
async def leave(self, room: MatrixRoom | str) -> None:
|
||||||
|
"""Leave the room (or reject the invite)."""
|
||||||
|
if self._config is None or self._client is None:
|
||||||
|
raise RuntimeError("ClientRoom is not set up")
|
||||||
|
res = await self._client.room_leave(self._room_id(room))
|
||||||
|
if not isinstance(res, RoomLeaveResponse):
|
||||||
|
raise RuntimeError(res)
|
||||||
|
|
||||||
|
async def join(self, room: MatrixRoom | str) -> None:
|
||||||
|
"""Join the room."""
|
||||||
|
if self._config is None or self._client is None:
|
||||||
|
raise RuntimeError("ClientRoom is not set up")
|
||||||
|
res = await self._client.join(
|
||||||
|
self._room_id(room)
|
||||||
|
)
|
||||||
|
if not isinstance(res, JoinResponse):
|
||||||
|
raise RuntimeError(res)
|
||||||
|
|
||||||
|
async def invite(self, room: MatrixRoom | str, user: MatrixUser | str) -> None:
|
||||||
|
"""Invite the user to the room."""
|
||||||
|
if self._config is None or self._client is None:
|
||||||
|
raise RuntimeError("ClientRoom is not set up")
|
||||||
|
res = await self._client.room_invite(
|
||||||
|
self._room_id(room),
|
||||||
|
self._user_id(user)
|
||||||
|
)
|
||||||
|
if not isinstance(res, RoomInviteResponse):
|
||||||
|
raise RuntimeError(res)
|
||||||
|
|
||||||
|
async def kick(self, room: MatrixRoom | str, user: MatrixUser | str, reason: str | None = None) -> None:
|
||||||
|
"""Kick the user from the room. The user will be able to join the room
|
||||||
|
again.
|
||||||
|
"""
|
||||||
|
if self._config is None or self._client is None:
|
||||||
|
raise RuntimeError("ClientRoom is not set up")
|
||||||
|
res = await self._client.room_kick(
|
||||||
|
self._room_id(room),
|
||||||
|
self._user_id(user),
|
||||||
|
reason
|
||||||
|
)
|
||||||
|
if not isinstance(res, RoomKickResponse):
|
||||||
|
raise RuntimeError(res)
|
||||||
|
|
||||||
|
async def ban(self, room: MatrixRoom | str, user: MatrixUser | str, reason: str | None = None) -> None:
|
||||||
|
"""Ban the user from the room. The user will be unable to join until he
|
||||||
|
is unbanned.
|
||||||
|
"""
|
||||||
|
if self._config is None or self._client is None:
|
||||||
|
raise RuntimeError("ClientRoom is not set up")
|
||||||
|
res = await self._client.room_ban(
|
||||||
|
self._room_id(room),
|
||||||
|
self._user_id(user),
|
||||||
|
reason
|
||||||
|
)
|
||||||
|
if not isinstance(res, RoomBanResponse):
|
||||||
|
raise RuntimeError(res)
|
||||||
|
|
||||||
|
async def unban(self, room: MatrixRoom | str, user: MatrixUser | str) -> None:
|
||||||
|
"""Unban the user from the room."""
|
||||||
|
if self._config is None or self._client is None:
|
||||||
|
raise RuntimeError("ClientRoom is not set up")
|
||||||
|
res = await self._client.room_unban(
|
||||||
|
self._room_id(room),
|
||||||
|
self._user_id(user)
|
||||||
|
)
|
||||||
|
if not isinstance(res, RoomUnbanResponse):
|
||||||
|
raise RuntimeError(res)
|
||||||
|
|
||||||
|
async def members(self, room: MatrixRoom | str) -> list[RoomMember]:
|
||||||
|
"""Get list of room members."""
|
||||||
|
if self._config is None or self._client is None:
|
||||||
|
raise RuntimeError("ClientRoom is not set up")
|
||||||
|
res = await self._client.joined_members(self._room_id(room))
|
||||||
|
if not isinstance(res, JoinedMembersResponse):
|
||||||
|
raise RuntimeError(res)
|
||||||
|
return res.members
|
||||||
@@ -98,7 +98,7 @@ class ClientSender:
|
|||||||
self._client = client
|
self._client = client
|
||||||
self._uploader = uploader
|
self._uploader = uploader
|
||||||
|
|
||||||
async def send_content(self,
|
async def content(self,
|
||||||
room: MatrixRoom | str,
|
room: MatrixRoom | str,
|
||||||
content: dict) -> RoomSendResponse:
|
content: dict) -> RoomSendResponse:
|
||||||
"""
|
"""
|
||||||
@@ -124,7 +124,7 @@ class ClientSender:
|
|||||||
if self._config.auto_verify_all_known_devices:
|
if self._config.auto_verify_all_known_devices:
|
||||||
if not Utils.verify_all_known_devices(self._client):
|
if not Utils.verify_all_known_devices(self._client):
|
||||||
raise
|
raise
|
||||||
return await self.send_content(room, content)
|
return await self.content(room, content)
|
||||||
else:
|
else:
|
||||||
raise
|
raise
|
||||||
if type(result) is RoomSendResponse:
|
if type(result) is RoomSendResponse:
|
||||||
@@ -134,7 +134,7 @@ class ClientSender:
|
|||||||
else:
|
else:
|
||||||
raise RuntimeError("Unknown error has occured", result)
|
raise RuntimeError("Unknown error has occured", result)
|
||||||
|
|
||||||
async def send_text(self,
|
async def text(self,
|
||||||
room: MatrixRoom | str,
|
room: MatrixRoom | str,
|
||||||
text: str,
|
text: str,
|
||||||
*,
|
*,
|
||||||
@@ -156,9 +156,9 @@ class ClientSender:
|
|||||||
"msgtype": "m.text",
|
"msgtype": "m.text",
|
||||||
**text_data
|
**text_data
|
||||||
}
|
}
|
||||||
return (await self.send_content(room, content)).event_id
|
return (await self.content(room, content)).event_id
|
||||||
|
|
||||||
async def send_image(self,
|
async def image(self,
|
||||||
room: MatrixRoom | str,
|
room: MatrixRoom | str,
|
||||||
path: Path | str, *,
|
path: Path | str, *,
|
||||||
text: str | None = None,
|
text: str | None = None,
|
||||||
@@ -213,9 +213,9 @@ class ClientSender:
|
|||||||
"h": height
|
"h": height
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return (await self.send_content(room, content)).event_id
|
return (await self.content(room, content)).event_id
|
||||||
|
|
||||||
async def send_image_bytes(self,
|
async def image_bytes(self,
|
||||||
room: MatrixRoom | str,
|
room: MatrixRoom | str,
|
||||||
data: bytes,
|
data: bytes,
|
||||||
filename: str,
|
filename: str,
|
||||||
@@ -276,9 +276,9 @@ class ClientSender:
|
|||||||
"h": height
|
"h": height
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return (await self.send_content(room, content)).event_id
|
return (await self.content(room, content)).event_id
|
||||||
|
|
||||||
async def send_video(self,
|
async def video(self,
|
||||||
room: MatrixRoom | str,
|
room: MatrixRoom | str,
|
||||||
path: Path | str,
|
path: Path | str,
|
||||||
*,
|
*,
|
||||||
@@ -361,4 +361,123 @@ class ClientSender:
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
# send
|
# send
|
||||||
return (await self.send_content(room, content)).event_id
|
return (await self.content(room, content)).event_id
|
||||||
|
|
||||||
|
async def file(self,
|
||||||
|
room: MatrixRoom | str,
|
||||||
|
path: Path | str,
|
||||||
|
*,
|
||||||
|
filename: str | None = None,
|
||||||
|
mime_type: str | None = None,
|
||||||
|
text: str | None = None,
|
||||||
|
is_html: bool | None = None,
|
||||||
|
timeout: float | None = 60 * 60):
|
||||||
|
"""
|
||||||
|
Send the file to `room`. Please note that formatted text is displayed
|
||||||
|
incorrectly in some clients as of September 8th, 2026.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
- room - the room to send the file to
|
||||||
|
- path - path to the file
|
||||||
|
- filename - filename to use for upload (`None` to use basename from
|
||||||
|
`path`)
|
||||||
|
- mime_type - mime type to use (`None` for autodetect using content
|
||||||
|
of `path`)
|
||||||
|
- text - caption to use (`None` to disable)
|
||||||
|
- is_html - whether the text is HTML-formatted (`None` for auto)
|
||||||
|
- timeout - upload timeout in seconds (`None` to disable)
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
- `event_id` of sent message on success
|
||||||
|
- Raises an exception on error
|
||||||
|
"""
|
||||||
|
if self._client is None or self._config is None:
|
||||||
|
raise RuntimeError("ClientSender is not set up")
|
||||||
|
# get filename if not set
|
||||||
|
if not filename:
|
||||||
|
filename = os.path.basename(path)
|
||||||
|
# caption must not actually be empty
|
||||||
|
if text is None or not text.strip():
|
||||||
|
text = filename
|
||||||
|
is_html = False
|
||||||
|
# get the mime type if not specified
|
||||||
|
if mime_type is None:
|
||||||
|
mime_type = magic.from_file(path, mime=True)
|
||||||
|
# upload
|
||||||
|
async with asyncio.timeout(timeout):
|
||||||
|
upload_result = await self._uploader.upload_file(
|
||||||
|
path, mime_type=mime_type, filename=filename)
|
||||||
|
# send
|
||||||
|
content = {
|
||||||
|
"msgtype": "m.file",
|
||||||
|
"filename": filename,
|
||||||
|
**self._process_html_text(text, is_html),
|
||||||
|
"file": {
|
||||||
|
"url": upload_result.response.content_uri,
|
||||||
|
"mimetype": mime_type,
|
||||||
|
**upload_result.keys
|
||||||
|
},
|
||||||
|
"info": {
|
||||||
|
"mimetype": mime_type,
|
||||||
|
"size": upload_result.filesize
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return (await self.content(room, content)).event_id
|
||||||
|
|
||||||
|
async def file_bytes(self,
|
||||||
|
room: MatrixRoom | str,
|
||||||
|
data: bytes,
|
||||||
|
filename: str,
|
||||||
|
*,
|
||||||
|
mime_type: str | None = None,
|
||||||
|
text: str | None = None,
|
||||||
|
is_html: bool | None = None,
|
||||||
|
timeout: float | None = 60 * 60) -> str:
|
||||||
|
"""
|
||||||
|
Send the file to `room`.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
- room - the room to send the file to
|
||||||
|
- data - the content of the file to send
|
||||||
|
- filename - filename to use for the file
|
||||||
|
- mime_type - mime type to use (`None` for autodetect using `data`
|
||||||
|
content)
|
||||||
|
- text - caption to use (`None` to disable)
|
||||||
|
- is_html - whether the text is HTML-formatted (`None` for auto)
|
||||||
|
- timeout - upload timeout in seconds (`None` to disable)
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
- `event_id` of sent message on success
|
||||||
|
- Raises an exception on error
|
||||||
|
"""
|
||||||
|
# caption must not actually be empty
|
||||||
|
if text is None or not text.strip():
|
||||||
|
text = filename
|
||||||
|
is_html = False
|
||||||
|
# check if the file is image
|
||||||
|
if not mime_type:
|
||||||
|
mime_type = magic.from_buffer(data, mime=True)
|
||||||
|
# upload
|
||||||
|
buffer = BytesIO(data)
|
||||||
|
async with asyncio.timeout(timeout):
|
||||||
|
upload_result = await self._uploader.upload_using_provider(
|
||||||
|
provider=buffer,
|
||||||
|
mime_type=mime_type,
|
||||||
|
filename=filename,
|
||||||
|
filesize=len(data))
|
||||||
|
# prepare the content and send
|
||||||
|
content = {
|
||||||
|
"msgtype": "m.file",
|
||||||
|
"filename": filename,
|
||||||
|
**self._process_html_text(text, is_html),
|
||||||
|
"file": {
|
||||||
|
"url": upload_result.response.content_uri,
|
||||||
|
"mimetype": mime_type,
|
||||||
|
**upload_result.keys
|
||||||
|
},
|
||||||
|
"info": {
|
||||||
|
"mimetype": mime_type,
|
||||||
|
"size": upload_result.filesize
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return (await self.content(room, content)).event_id
|
||||||
@@ -1,16 +1,19 @@
|
|||||||
import asyncio
|
import asyncio
|
||||||
import logging
|
import logging
|
||||||
|
|
||||||
from typing import Callable, Coroutine, Any
|
from typing import Callable, Coroutine, Any, overload
|
||||||
|
|
||||||
from nio import AsyncClient
|
from nio import AsyncClient, MatrixRoom, Event
|
||||||
|
|
||||||
from ..filters.base import BaseEventFilter
|
from ..filters.base import BaseEventFilter
|
||||||
from ..types import *
|
from ..types import *
|
||||||
|
from ..context import EventContext
|
||||||
|
|
||||||
from ._validation import Validator
|
from ._validation import Validator
|
||||||
from ._storage import Storage
|
from ._storage import Storage
|
||||||
from ._client_auth import ClientAuth
|
from ._client_auth import ClientAuth
|
||||||
|
from ._client_room_operations import ClientRoomOperations
|
||||||
|
from ._client_downloader import ClientDownloader
|
||||||
from ._client_manager import ClientManager
|
from ._client_manager import ClientManager
|
||||||
from ._client_uploader import ClientUploader
|
from ._client_uploader import ClientUploader
|
||||||
from ._client_sender import ClientSender
|
from ._client_sender import ClientSender
|
||||||
@@ -32,6 +35,8 @@ class MatrixBot:
|
|||||||
self._client_auth = ClientAuth(self._storage)
|
self._client_auth = ClientAuth(self._storage)
|
||||||
self._client_manager = ClientManager(self._client_auth, self._storage)
|
self._client_manager = ClientManager(self._client_auth, self._storage)
|
||||||
self._client_uploader = ClientUploader(self._storage)
|
self._client_uploader = ClientUploader(self._storage)
|
||||||
|
self._client_downloader = ClientDownloader()
|
||||||
|
self._client_room_operations = ClientRoomOperations()
|
||||||
self._client_sender = ClientSender()
|
self._client_sender = ClientSender()
|
||||||
self._callbacks = Callbacks(self._storage, self)
|
self._callbacks = Callbacks(self._storage, self)
|
||||||
# validate the config and save it
|
# validate the config and save it
|
||||||
@@ -44,7 +49,7 @@ class MatrixBot:
|
|||||||
|
|
||||||
def add_callback(self,
|
def add_callback(self,
|
||||||
filter: BaseEventFilter,
|
filter: BaseEventFilter,
|
||||||
callback: Callable[[RoomEventData], Coroutine[Any, Any, None]] | None,
|
callback: Callable[[EventContext], Coroutine[Any, Any, None]] | None,
|
||||||
*,
|
*,
|
||||||
stop_matching: bool = True) -> None:
|
stop_matching: bool = True) -> None:
|
||||||
"""
|
"""
|
||||||
@@ -79,6 +84,14 @@ class MatrixBot:
|
|||||||
self._config,
|
self._config,
|
||||||
self._client_manager.get_client()
|
self._client_manager.get_client()
|
||||||
)
|
)
|
||||||
|
await self._client_downloader.setup(
|
||||||
|
self._config,
|
||||||
|
self._client_manager.get_client()
|
||||||
|
)
|
||||||
|
await self._client_room_operations.setup(
|
||||||
|
self._config,
|
||||||
|
self._client_manager.get_client()
|
||||||
|
)
|
||||||
await self._client_sender.setup(
|
await self._client_sender.setup(
|
||||||
self._config,
|
self._config,
|
||||||
self._client_manager.get_client(),
|
self._client_manager.get_client(),
|
||||||
@@ -112,7 +125,8 @@ class MatrixBot:
|
|||||||
finally:
|
finally:
|
||||||
await self.stop()
|
await self.stop()
|
||||||
|
|
||||||
def get_client(self) -> AsyncClient:
|
@property
|
||||||
|
def client(self) -> AsyncClient:
|
||||||
"""
|
"""
|
||||||
Get AsyncClient.
|
Get AsyncClient.
|
||||||
|
|
||||||
@@ -122,125 +136,19 @@ class MatrixBot:
|
|||||||
"""
|
"""
|
||||||
return self._client_manager.get_client()
|
return self._client_manager.get_client()
|
||||||
|
|
||||||
async def send_text(self,
|
@property
|
||||||
room: MatrixRoom | str,
|
def send(self) -> ClientSender:
|
||||||
text: str,
|
"""Get ClientSender that you should use to send messages."""
|
||||||
*,
|
return self._client_sender
|
||||||
is_html: bool | None = None) -> str:
|
|
||||||
|
@property
|
||||||
|
def download(self) -> ClientDownloader:
|
||||||
|
"""Get ClientDownloader that you should use to download files."""
|
||||||
|
return self._client_downloader
|
||||||
|
|
||||||
|
@property
|
||||||
|
def rooms(self) -> ClientRoomOperations:
|
||||||
|
"""Get ClientRoomOperations that you should use to do room-related
|
||||||
|
actions.
|
||||||
"""
|
"""
|
||||||
Send text message to `room`.
|
return self._client_room_operations
|
||||||
|
|
||||||
Args:
|
|
||||||
- room - the room to send the text to
|
|
||||||
- text - the text to send to the room
|
|
||||||
- is_html - whether the text is HTML-formatted. Use `None` for auto
|
|
||||||
|
|
||||||
Returns:
|
|
||||||
- `event_id` of sent message on success
|
|
||||||
- Raises an exception on error
|
|
||||||
"""
|
|
||||||
return await self._client_sender.send_text(
|
|
||||||
room=room,
|
|
||||||
text=text,
|
|
||||||
is_html=is_html
|
|
||||||
)
|
|
||||||
|
|
||||||
async def send_image(self,
|
|
||||||
room: MatrixRoom | str,
|
|
||||||
path: Path | str, *,
|
|
||||||
text: str | None = None,
|
|
||||||
is_html: bool | None = None,
|
|
||||||
filename: str | None = None,
|
|
||||||
timeout: float | None = 60 * 60) -> str:
|
|
||||||
"""
|
|
||||||
Send the image to `room`. Please note that formatted text is displayed
|
|
||||||
incorrectly in some clients as of September 8th, 2026
|
|
||||||
|
|
||||||
Args:
|
|
||||||
- room - the room to send the text to
|
|
||||||
- path - path to the image file
|
|
||||||
- text - image caption to use (`None` to disable)
|
|
||||||
- is_html - whether the text is HTML-formatted (`None` for auto)
|
|
||||||
- filename - filename to use for the file (`None` for auto)
|
|
||||||
- timeout - upload timeout in seconds (`None` to disable)
|
|
||||||
|
|
||||||
Returns:
|
|
||||||
- `event_id` of sent message on success
|
|
||||||
- Raises an exception on error
|
|
||||||
"""
|
|
||||||
return await self._client_sender.send_image(
|
|
||||||
room=room,
|
|
||||||
path=path,
|
|
||||||
text=text,
|
|
||||||
is_html=is_html,
|
|
||||||
filename=filename,
|
|
||||||
timeout=timeout
|
|
||||||
)
|
|
||||||
|
|
||||||
async def send_image_bytes(self,
|
|
||||||
room: MatrixRoom | str,
|
|
||||||
data: bytes,
|
|
||||||
filename: str,
|
|
||||||
*,
|
|
||||||
text: str | None = None,
|
|
||||||
is_html: bool | None = None,
|
|
||||||
timeout: float | None = 60 * 60) -> str:
|
|
||||||
"""
|
|
||||||
Send the image to `room`. Please note that formatted text is displayed
|
|
||||||
incorrectly in some clients as of September 8th, 2026
|
|
||||||
|
|
||||||
Args:
|
|
||||||
- room - the room to send the text to
|
|
||||||
- bytes - the image to send
|
|
||||||
- filename - filename to use for the file
|
|
||||||
- text - image caption to use (`None` to disable)
|
|
||||||
- is_html - whether the text is HTML-formatted (`None` for auto)
|
|
||||||
- timeout - upload timeout in seconds (`None` to disable)
|
|
||||||
|
|
||||||
Returns:
|
|
||||||
- `event_id` of sent message on success
|
|
||||||
- Raises an exception on error
|
|
||||||
"""
|
|
||||||
return await self._client_sender.send_image_bytes(
|
|
||||||
room=room,
|
|
||||||
data=data,
|
|
||||||
filename=filename,
|
|
||||||
text=text,
|
|
||||||
is_html=is_html,
|
|
||||||
timeout=timeout
|
|
||||||
)
|
|
||||||
|
|
||||||
async def send_video(self,
|
|
||||||
room: MatrixRoom | str,
|
|
||||||
path: Path | str,
|
|
||||||
*,
|
|
||||||
props: VideoFileProperties | None = None,
|
|
||||||
text: str | None = None,
|
|
||||||
is_html: bool | None = None,
|
|
||||||
timeout: float | None = 60 * 60) -> str:
|
|
||||||
"""
|
|
||||||
Send the video to `room`. Please note that formatted text is displayed
|
|
||||||
incorrectly in some clients as of September 8th, 2026. Unknown video
|
|
||||||
properties will be automatically deduced as configured in
|
|
||||||
`MatrixBotConfig`.
|
|
||||||
|
|
||||||
Args:
|
|
||||||
- room - the room to send the text to
|
|
||||||
- path - path to the video file
|
|
||||||
- props - video properties (`None` for auto, if the feature is ON)
|
|
||||||
- text - video caption to use (`None` to disable)
|
|
||||||
- is_html - whether the text is HTML-formatted (`None` for auto)
|
|
||||||
- timeout - upload timeout in seconds (`None` to disable)
|
|
||||||
|
|
||||||
Returns:
|
|
||||||
- `event_id` of sent message on success
|
|
||||||
- Raises an exception on error
|
|
||||||
"""
|
|
||||||
return await self._client_sender.send_video(
|
|
||||||
room=room,
|
|
||||||
path=path,
|
|
||||||
props=props,
|
|
||||||
text=text,
|
|
||||||
is_html=is_html,
|
|
||||||
timeout=timeout
|
|
||||||
)
|
|
||||||
88
src/mab/context.py
Normal file
88
src/mab/context.py
Normal file
@@ -0,0 +1,88 @@
|
|||||||
|
"""This module implements logic for event context"""
|
||||||
|
|
||||||
|
from dataclasses import dataclass
|
||||||
|
from typing import Any, TYPE_CHECKING
|
||||||
|
|
||||||
|
from .types import ContextDataKey, MessageType
|
||||||
|
|
||||||
|
from nio import MatrixRoom, Event
|
||||||
|
|
||||||
|
if TYPE_CHECKING:
|
||||||
|
from .bot import MatrixBot
|
||||||
|
|
||||||
|
#
|
||||||
|
# Possible context variables
|
||||||
|
#
|
||||||
|
CTX_BODY = ContextDataKey[str]("CTX_BODY")
|
||||||
|
"""Value of `event.body`"""
|
||||||
|
|
||||||
|
CTX_MESSAGE_TYPE = ContextDataKey[MessageType]("CTX_MESSAGE_TYPE")
|
||||||
|
"""Value of `msgtype` for the event"""
|
||||||
|
|
||||||
|
CTX_SENDER = ContextDataKey[str]("CTX_SENDER")
|
||||||
|
"""Value of `event.sender`"""
|
||||||
|
|
||||||
|
CTX_CMD_PREFIX = ContextDataKey[str]("CTX_CMD_PREFIX")
|
||||||
|
"""Command prefix that was used when matching"""
|
||||||
|
|
||||||
|
CTX_CMD_VERB = ContextDataKey[str]("CTX_CMD_VERB")
|
||||||
|
"""The verb that was used to execute the command"""
|
||||||
|
|
||||||
|
CTX_CMD_ARGS = ContextDataKey[list[str]]("CTX_CMD_ARGS")
|
||||||
|
"""Arguments that were passed with the command"""
|
||||||
|
|
||||||
|
CTX_ROOM_ENCRYPTED = ContextDataKey[bool]("CTX_ROOM_ENCRYPTED")
|
||||||
|
"""True if the room is encrypted"""
|
||||||
|
|
||||||
|
CTX_FILE_SIZE = ContextDataKey[int]("CTX_FILE_SIZE")
|
||||||
|
"""Size of the file attached to the message (bytes)"""
|
||||||
|
|
||||||
|
CTX_FILE_MIME = ContextDataKey[str]("CTX_FILE_MIME")
|
||||||
|
"""Mime type of the file attached to the message"""
|
||||||
|
|
||||||
|
CTX_FILE_NAME = ContextDataKey[str]("CTX_FILE_NAME")
|
||||||
|
"""Name of the file attached to the message"""
|
||||||
|
|
||||||
|
|
||||||
|
#
|
||||||
|
# EventContext implementation
|
||||||
|
#
|
||||||
|
@dataclass
|
||||||
|
class EventContext:
|
||||||
|
"""The class holding information about an event that happened in the room"""
|
||||||
|
|
||||||
|
room: MatrixRoom
|
||||||
|
"""The room the event has happened in"""
|
||||||
|
|
||||||
|
event: Event
|
||||||
|
"""The event that has happened in the room"""
|
||||||
|
|
||||||
|
bot: "MatrixBot"
|
||||||
|
"""The bot that is the source of the event"""
|
||||||
|
|
||||||
|
def __setitem__[T](self, key: ContextDataKey[T], value: T | None) -> None:
|
||||||
|
"""Set a value inside the context data storage. `None` removes it"""
|
||||||
|
if not hasattr(self, "_datastore"):
|
||||||
|
self._datastore: dict[ContextDataKey, Any] = {}
|
||||||
|
if value is None:
|
||||||
|
del self._datastore[key]
|
||||||
|
else:
|
||||||
|
self._datastore[key] = value
|
||||||
|
|
||||||
|
def __getitem__[T](self, key: ContextDataKey[T]) -> T:
|
||||||
|
"""
|
||||||
|
Get a value inside the context data storage.
|
||||||
|
|
||||||
|
Raises RuntimeError if the value is not present.
|
||||||
|
"""
|
||||||
|
if not hasattr(self, "_datastore") or key not in self._datastore:
|
||||||
|
raise RuntimeError(f"Context does not contain {repr(key)}")
|
||||||
|
return self._datastore[key]
|
||||||
|
|
||||||
|
def __contains__[T](self, key: ContextDataKey[T]) -> bool:
|
||||||
|
"""Check if context data storage contains the value"""
|
||||||
|
if not hasattr(self, "_datastore"):
|
||||||
|
return False
|
||||||
|
if key not in self._datastore:
|
||||||
|
return False
|
||||||
|
return True
|
||||||
@@ -5,6 +5,8 @@ from typing import Any, Type
|
|||||||
from nio import AsyncClient
|
from nio import AsyncClient
|
||||||
from nio import MatrixRoom, Event
|
from nio import MatrixRoom, Event
|
||||||
|
|
||||||
|
from ..context import EventContext
|
||||||
|
|
||||||
class BaseEventFilter(ABC):
|
class BaseEventFilter(ABC):
|
||||||
"""Base class for all message filters"""
|
"""Base class for all message filters"""
|
||||||
_logger = logging.Logger("EventFilter")
|
_logger = logging.Logger("EventFilter")
|
||||||
@@ -35,7 +37,7 @@ class BaseEventFilter(ABC):
|
|||||||
)
|
)
|
||||||
|
|
||||||
def __ror__(self, other):
|
def __ror__(self, other):
|
||||||
return self.__ror__(other)
|
return self.__or__(other)
|
||||||
|
|
||||||
# XOR
|
# XOR
|
||||||
def __xor__(self, other):
|
def __xor__(self, other):
|
||||||
@@ -64,7 +66,7 @@ class BaseEventFilter(ABC):
|
|||||||
"""
|
"""
|
||||||
return str(self.__class__.__name__)
|
return str(self.__class__.__name__)
|
||||||
|
|
||||||
async def __call__(self, room: MatrixRoom, event: Event, client: AsyncClient) -> bool:
|
async def __call__(self, context: EventContext) -> bool:
|
||||||
"""
|
"""
|
||||||
This abstract method must be redefined in derived classes so that the
|
This abstract method must be redefined in derived classes so that the
|
||||||
filter operates according to its description. This method must not raise
|
filter operates according to its description. This method must not raise
|
||||||
@@ -72,9 +74,9 @@ class BaseEventFilter(ABC):
|
|||||||
and return False
|
and return False
|
||||||
|
|
||||||
Args:
|
Args:
|
||||||
- room - room the event has happened in
|
- context - event context; your derived classes may add variables
|
||||||
- event - the event to check againts this filter
|
to it (see `message.MessageTypeFilter` implementation
|
||||||
- client - the client
|
for reference)
|
||||||
|
|
||||||
Returns:
|
Returns:
|
||||||
- True if the event satisfies this filter
|
- True if the event satisfies this filter
|
||||||
@@ -88,8 +90,10 @@ class EventTypeFilter(BaseEventFilter):
|
|||||||
super().__init__(**kwargs)
|
super().__init__(**kwargs)
|
||||||
self._type = type
|
self._type = type
|
||||||
|
|
||||||
async def __call__(self, room: MatrixRoom, event: Event, client: AsyncClient) -> bool:
|
async def __call__(self, context: EventContext) -> bool:
|
||||||
return isinstance(event, self._type)
|
if not await super().__call__(context):
|
||||||
|
return False
|
||||||
|
return isinstance(context.event, self._type)
|
||||||
|
|
||||||
class CompoundEventFilter(BaseEventFilter):
|
class CompoundEventFilter(BaseEventFilter):
|
||||||
"""Event filter that consists of multiple filters"""
|
"""Event filter that consists of multiple filters"""
|
||||||
@@ -140,8 +144,10 @@ class CompoundEventFilter(BaseEventFilter):
|
|||||||
expression = f"~{reprs[0]}"
|
expression = f"~{reprs[0]}"
|
||||||
return f"({expression})"
|
return f"({expression})"
|
||||||
|
|
||||||
async def __call__(self, room: MatrixRoom, event: Event, client: AsyncClient) -> bool:
|
async def __call__(self, context: EventContext) -> bool:
|
||||||
evaluated = [await arg(room, event, client) for arg in self._arguments]
|
if not await super().__call__(context):
|
||||||
|
return False
|
||||||
|
evaluated = [await arg(context) for arg in self._arguments]
|
||||||
if self._operator == self.OPERATOR_AND:
|
if self._operator == self.OPERATOR_AND:
|
||||||
return all(evaluated)
|
return all(evaluated)
|
||||||
elif self._operator == self.OPERATOR_OR:
|
elif self._operator == self.OPERATOR_OR:
|
||||||
|
|||||||
@@ -1,6 +1,13 @@
|
|||||||
import re
|
import re
|
||||||
import traceback
|
import traceback
|
||||||
from .message import NewMessageFilter
|
from .message import NewMessageFilter
|
||||||
|
from ..context import EventContext
|
||||||
|
from ..context import (
|
||||||
|
CTX_BODY,
|
||||||
|
CTX_CMD_PREFIX,
|
||||||
|
CTX_CMD_VERB,
|
||||||
|
CTX_CMD_ARGS
|
||||||
|
)
|
||||||
|
|
||||||
from nio import AsyncClient
|
from nio import AsyncClient
|
||||||
from nio import MatrixRoom, Event
|
from nio import MatrixRoom, Event
|
||||||
@@ -21,24 +28,27 @@ class BodyExistsFilter(NewMessageFilter):
|
|||||||
|
|
||||||
This filter will match any message that has `body` in it, including images,
|
This filter will match any message that has `body` in it, including images,
|
||||||
videos, files, etc.
|
videos, files, etc.
|
||||||
|
|
||||||
|
This filter sets `CTX_BODY` context variable.
|
||||||
"""
|
"""
|
||||||
def __init__(self, *, ignore_filename_in_body: bool = True, **kwargs):
|
def __init__(self, *, ignore_filename_in_body: bool = True, **kwargs):
|
||||||
super().__init__(**kwargs)
|
super().__init__(**kwargs)
|
||||||
self._ignore_filename_in_body = ignore_filename_in_body
|
self._ignore_filename_in_body = ignore_filename_in_body
|
||||||
|
|
||||||
async def __call__(self, room: MatrixRoom, event: Event, client: AsyncClient) -> bool:
|
async def __call__(self, context: EventContext) -> bool:
|
||||||
if not await super().__call__(room, event, client):
|
if not await super().__call__(context):
|
||||||
return False
|
return False
|
||||||
if not hasattr(event, "body"):
|
if not hasattr(context.event, "body"):
|
||||||
return False
|
return False
|
||||||
if not isinstance(event.body, str): # type: ignore
|
if not isinstance(context.event.body, str): # type: ignore
|
||||||
return False
|
return False
|
||||||
if not event.body.strip(): # type: ignore
|
if not context.event.body.strip(): # type: ignore
|
||||||
return False
|
return False
|
||||||
if self._ignore_filename_in_body:
|
if self._ignore_filename_in_body:
|
||||||
content = event.source["content"]
|
content = context.event.source["content"]
|
||||||
if "filename" in content and content["filename"] == event.body: # type: ignore
|
if "filename" in content and content["filename"] == context.event.body: # type: ignore
|
||||||
return False
|
return False
|
||||||
|
context[CTX_BODY] = context.event.body # type: ignore
|
||||||
return True
|
return True
|
||||||
|
|
||||||
class BodyContainsFilter(BodyExistsFilter):
|
class BodyContainsFilter(BodyExistsFilter):
|
||||||
@@ -61,10 +71,10 @@ class BodyContainsFilter(BodyExistsFilter):
|
|||||||
self._any_case = any_case
|
self._any_case = any_case
|
||||||
self._needle = needle
|
self._needle = needle
|
||||||
|
|
||||||
async def __call__(self, room: MatrixRoom, event: Event, client: AsyncClient) -> bool:
|
async def __call__(self, context: EventContext) -> bool:
|
||||||
if not await super().__call__(room, event, client):
|
if not await super().__call__(context):
|
||||||
return False
|
return False
|
||||||
body = event.body.lower() if self._any_case else event.body # type: ignore
|
body = context.event.body.lower() if self._any_case else context.event.body # type: ignore
|
||||||
for n in self._needle:
|
for n in self._needle:
|
||||||
if n in body:
|
if n in body:
|
||||||
return True
|
return True
|
||||||
@@ -90,10 +100,10 @@ class BodyStartsWithFilter(BodyExistsFilter):
|
|||||||
self._any_case = any_case
|
self._any_case = any_case
|
||||||
self._substring = substring
|
self._substring = substring
|
||||||
|
|
||||||
async def __call__(self, room: MatrixRoom, event: Event, client: AsyncClient) -> bool:
|
async def __call__(self, context: EventContext) -> bool:
|
||||||
if not await super().__call__(room, event, client):
|
if not await super().__call__(context):
|
||||||
return False
|
return False
|
||||||
body = event.body.lower() if self._any_case else event.body # type: ignore
|
body = context.event.body.lower() if self._any_case else context.event.body # type: ignore
|
||||||
for s in self._substring:
|
for s in self._substring:
|
||||||
if body.startswith(s):
|
if body.startswith(s):
|
||||||
return True
|
return True
|
||||||
@@ -119,10 +129,10 @@ class BodyEndsWithFilter(BodyExistsFilter):
|
|||||||
self._any_case = any_case
|
self._any_case = any_case
|
||||||
self._substring = substring
|
self._substring = substring
|
||||||
|
|
||||||
async def __call__(self, room: MatrixRoom, event: Event, client: AsyncClient) -> bool:
|
async def __call__(self, context: EventContext) -> bool:
|
||||||
if not await super().__call__(room, event, client):
|
if not await super().__call__(context):
|
||||||
return False
|
return False
|
||||||
body = event.body.lower() if self._any_case else event.body # type: ignore
|
body = context.event.body.lower() if self._any_case else context.event.body # type: ignore
|
||||||
for s in self._substring:
|
for s in self._substring:
|
||||||
if body.endswith(s):
|
if body.endswith(s):
|
||||||
return True
|
return True
|
||||||
@@ -145,9 +155,10 @@ class BodyCommandFilter(BodyExistsFilter):
|
|||||||
store all verbs in lower case. This filter will not match any verbs that
|
store all verbs in lower case. This filter will not match any verbs that
|
||||||
use mixed case of upper case.
|
use mixed case of upper case.
|
||||||
|
|
||||||
If this filter is matched, then it will set a new attribute for the event:
|
This filter sets the following context variables:
|
||||||
`event.command_args: list[str]`. You may use this attribute in your callback
|
- `CTX_CMD_PREFIX` - prefix that was used
|
||||||
for this event.
|
- `CTX_CMD_VERB` - verb that was used
|
||||||
|
- `CTX_CMD_ARGS` - arguments that were passed
|
||||||
"""
|
"""
|
||||||
def __init__(self, verbs: str | list[str], min_args: int = 0, max_args: int | None = None, prefix: str = "!", **kwargs):
|
def __init__(self, verbs: str | list[str], min_args: int = 0, max_args: int | None = None, prefix: str = "!", **kwargs):
|
||||||
super().__init__(**kwargs)
|
super().__init__(**kwargs)
|
||||||
@@ -158,10 +169,10 @@ class BodyCommandFilter(BodyExistsFilter):
|
|||||||
self._max_args = max_args
|
self._max_args = max_args
|
||||||
self._prefix = prefix
|
self._prefix = prefix
|
||||||
|
|
||||||
async def __call__(self, room: MatrixRoom, event: Event, client: AsyncClient) -> bool:
|
async def __call__(self, context: EventContext) -> bool:
|
||||||
if not await super().__call__(room, event, client):
|
if not await super().__call__(context):
|
||||||
return False
|
return False
|
||||||
parts = [p.strip() for p in event.body.split() if p.strip()] # type: ignore
|
parts = [p.strip() for p in context.event.body.split() if p.strip()] # type: ignore
|
||||||
args_count = len(parts) - 1
|
args_count = len(parts) - 1
|
||||||
if args_count < self._min_args:
|
if args_count < self._min_args:
|
||||||
return False
|
return False
|
||||||
@@ -172,7 +183,9 @@ class BodyCommandFilter(BodyExistsFilter):
|
|||||||
cmd = parts[0][len(self._prefix):].lower()
|
cmd = parts[0][len(self._prefix):].lower()
|
||||||
for verb in self._verbs:
|
for verb in self._verbs:
|
||||||
if cmd == verb:
|
if cmd == verb:
|
||||||
setattr(event, "command_args", parts[1:])
|
context[CTX_CMD_PREFIX] = self._prefix
|
||||||
|
context[CTX_CMD_VERB] = verb
|
||||||
|
context[CTX_CMD_ARGS] = parts[1:]
|
||||||
return True
|
return True
|
||||||
return False
|
return False
|
||||||
|
|
||||||
@@ -186,11 +199,11 @@ class BodyRegexFilter(BodyExistsFilter):
|
|||||||
regex = re.compile(regex)
|
regex = re.compile(regex)
|
||||||
self._regex = regex
|
self._regex = regex
|
||||||
|
|
||||||
async def __call__(self, room: MatrixRoom, event: Event, client: AsyncClient) -> bool:
|
async def __call__(self, context: EventContext) -> bool:
|
||||||
if not await super().__call__(room, event, client):
|
if not await super().__call__(context):
|
||||||
return False
|
return False
|
||||||
try:
|
try:
|
||||||
return self._regex.match(event.body) # type: ignore
|
return self._regex.match(context.event.body) is not None # type: ignore
|
||||||
except:
|
except:
|
||||||
self._logger.error(traceback.format_exc())
|
self._logger.error(traceback.format_exc())
|
||||||
return False
|
return False
|
||||||
@@ -1,6 +1,12 @@
|
|||||||
import traceback
|
import traceback
|
||||||
from .base import BaseEventFilter, EventTypeFilter
|
from .base import BaseEventFilter, EventTypeFilter
|
||||||
from ..types import MessageType
|
from ..types import MessageType
|
||||||
|
from ..context import (EventContext,
|
||||||
|
CTX_MESSAGE_TYPE,
|
||||||
|
CTX_SENDER,
|
||||||
|
CTX_FILE_SIZE,
|
||||||
|
CTX_FILE_MIME,
|
||||||
|
CTX_FILE_NAME)
|
||||||
|
|
||||||
from nio import AsyncClient
|
from nio import AsyncClient
|
||||||
from nio import MatrixRoom, Event
|
from nio import MatrixRoom, Event
|
||||||
@@ -15,6 +21,8 @@ class MessageTypeFilter(BaseEventFilter):
|
|||||||
|
|
||||||
`types` list is stored by reference so you may modify the behavior of this
|
`types` list is stored by reference so you may modify the behavior of this
|
||||||
filter dynamically.
|
filter dynamically.
|
||||||
|
|
||||||
|
This filter sets `CTX_MESSAGE_TYPE` variable in the context.
|
||||||
"""
|
"""
|
||||||
def __init__(self, types: list[MessageType] | MessageType, **kwargs):
|
def __init__(self, types: list[MessageType] | MessageType, **kwargs):
|
||||||
super().__init__(**kwargs)
|
super().__init__(**kwargs)
|
||||||
@@ -22,14 +30,16 @@ class MessageTypeFilter(BaseEventFilter):
|
|||||||
types = [types]
|
types = [types]
|
||||||
self._types = types
|
self._types = types
|
||||||
|
|
||||||
async def __call__(self, room: MatrixRoom, event: Event, client: AsyncClient) -> bool:
|
async def __call__(self, context: EventContext) -> bool:
|
||||||
if not await super().__call__(room, event, client):
|
if not await super().__call__(context):
|
||||||
return False
|
return False
|
||||||
if "msgtype" not in event.source["content"]:
|
if "msgtype" not in context.event.source["content"]:
|
||||||
return False
|
return False
|
||||||
return (
|
msgtype = context.event.source["content"]["msgtype"]
|
||||||
event.source["content"]["msgtype"] in [t.value for t in self._types]
|
if not msgtype in [t.value for t in self._types]:
|
||||||
)
|
return False
|
||||||
|
context[CTX_MESSAGE_TYPE] = MessageType(msgtype)
|
||||||
|
return True
|
||||||
|
|
||||||
class NewMessageFilter(BaseEventFilter):
|
class NewMessageFilter(BaseEventFilter):
|
||||||
"""
|
"""
|
||||||
@@ -40,10 +50,10 @@ class NewMessageFilter(BaseEventFilter):
|
|||||||
def __init__(self, **kwargs):
|
def __init__(self, **kwargs):
|
||||||
super().__init__(**kwargs)
|
super().__init__(**kwargs)
|
||||||
|
|
||||||
async def __call__(self, room: MatrixRoom, event: Event, client: AsyncClient) -> bool:
|
async def __call__(self, context: EventContext) -> bool:
|
||||||
if not await super().__call__(room, event, client):
|
if not await super().__call__(context):
|
||||||
return False
|
return False
|
||||||
return "m.new_content" not in event.source["content"]
|
return "m.new_content" not in context.event.source["content"]
|
||||||
|
|
||||||
class EditedMessageFilter(BaseEventFilter):
|
class EditedMessageFilter(BaseEventFilter):
|
||||||
"""
|
"""
|
||||||
@@ -53,10 +63,10 @@ class EditedMessageFilter(BaseEventFilter):
|
|||||||
def __init__(self, **kwargs):
|
def __init__(self, **kwargs):
|
||||||
super().__init__(**kwargs)
|
super().__init__(**kwargs)
|
||||||
|
|
||||||
async def __call__(self, room: MatrixRoom, event: Event, client: AsyncClient) -> bool:
|
async def __call__(self, context: EventContext) -> bool:
|
||||||
if not await super().__call__(room, event, client):
|
if not await super().__call__(context):
|
||||||
return False
|
return False
|
||||||
return "m.new_content" in event.source["content"]
|
return "m.new_content" in context.event.source["content"]
|
||||||
|
|
||||||
class RedactedMessageFilter(EventTypeFilter):
|
class RedactedMessageFilter(EventTypeFilter):
|
||||||
"""
|
"""
|
||||||
@@ -64,9 +74,10 @@ class RedactedMessageFilter(EventTypeFilter):
|
|||||||
"""
|
"""
|
||||||
def __init__(self, **kwargs):
|
def __init__(self, **kwargs):
|
||||||
super().__init__(RedactionEvent, **kwargs)
|
super().__init__(RedactionEvent, **kwargs)
|
||||||
|
raise NotImplementedError()
|
||||||
|
|
||||||
async def __call__(self, room: MatrixRoom, event: Event, client: AsyncClient) -> bool:
|
async def __call__(self, context: EventContext) -> bool:
|
||||||
return await super().__call__(room, event, client)
|
raise NotImplementedError()
|
||||||
|
|
||||||
class SenderIsFilter(BaseEventFilter):
|
class SenderIsFilter(BaseEventFilter):
|
||||||
"""
|
"""
|
||||||
@@ -77,6 +88,8 @@ class SenderIsFilter(BaseEventFilter):
|
|||||||
|
|
||||||
`senders` list is stored by reference so you can modify behavior of this
|
`senders` list is stored by reference so you can modify behavior of this
|
||||||
filter dynamically.
|
filter dynamically.
|
||||||
|
|
||||||
|
This filter sets `CTX_SENDER` variable in the context.
|
||||||
"""
|
"""
|
||||||
def __init__(self, sender: list[str] | str, *, any_case: bool = True, **kwargs):
|
def __init__(self, sender: list[str] | str, *, any_case: bool = True, **kwargs):
|
||||||
super().__init__(**kwargs)
|
super().__init__(**kwargs)
|
||||||
@@ -85,12 +98,13 @@ class SenderIsFilter(BaseEventFilter):
|
|||||||
self._sender = sender
|
self._sender = sender
|
||||||
self._any_case = any_case
|
self._any_case = any_case
|
||||||
|
|
||||||
async def __call__(self, room: MatrixRoom, event: Event, client: AsyncClient) -> bool:
|
async def __call__(self, context: EventContext) -> bool:
|
||||||
if not await super().__call__(room, event, client):
|
if not await super().__call__(context):
|
||||||
return False
|
return False
|
||||||
sender = event.sender.lower() if self._any_case else event.sender
|
sender = context.event.sender.lower() if self._any_case else context.event.sender
|
||||||
for s in self._sender:
|
for s in self._sender:
|
||||||
if sender == s:
|
if sender == s:
|
||||||
|
context[CTX_SENDER] = sender
|
||||||
return True
|
return True
|
||||||
return False
|
return False
|
||||||
|
|
||||||
@@ -106,7 +120,32 @@ class SenderIsBotFilter(BaseEventFilter):
|
|||||||
def __init__(self, **kwargs):
|
def __init__(self, **kwargs):
|
||||||
super().__init__(**kwargs)
|
super().__init__(**kwargs)
|
||||||
|
|
||||||
async def __call__(self, room: MatrixRoom, event: Event, client: AsyncClient) -> bool:
|
async def __call__(self, context: EventContext) -> bool:
|
||||||
if not await super().__call__(room, event, client):
|
if not await super().__call__(context):
|
||||||
return False
|
return False
|
||||||
return client.user_id == event.sender
|
return context.bot.client.user_id == context.event.sender
|
||||||
|
|
||||||
|
class MessageHasFile(BaseEventFilter):
|
||||||
|
"""
|
||||||
|
This filter returns True if the message contains file that can be
|
||||||
|
downloaded.
|
||||||
|
|
||||||
|
This filter sets `CTX_FILE_NAME`, `CTX_FILE_SIZE` and `CTX_FILE_MIME`
|
||||||
|
variables in the context (if they are present in the).
|
||||||
|
"""
|
||||||
|
def __init__(self, **kwargs):
|
||||||
|
super().__init__(**kwargs)
|
||||||
|
|
||||||
|
async def __call__(self, context: EventContext) -> bool:
|
||||||
|
if not await super().__call__(context):
|
||||||
|
return False
|
||||||
|
content: dict | None = context.event.source.get("content")
|
||||||
|
if content is None:
|
||||||
|
return False
|
||||||
|
# `file` if encrypted, `url` if not encrypted
|
||||||
|
if "file" not in content and "url" not in content:
|
||||||
|
return False
|
||||||
|
context[CTX_FILE_NAME] = content.get("filename") or content.get("body")
|
||||||
|
context[CTX_FILE_SIZE] = content["info"].get("size")
|
||||||
|
context[CTX_FILE_MIME] = content["info"].get("mimetype")
|
||||||
|
return True
|
||||||
@@ -1,7 +1,6 @@
|
|||||||
from .base import BaseEventFilter
|
from .base import BaseEventFilter
|
||||||
|
|
||||||
from nio import AsyncClient
|
from ..context import EventContext, CTX_ROOM_ENCRYPTED
|
||||||
from nio import MatrixRoom, Event
|
|
||||||
|
|
||||||
class RoomEncryptedFilter(BaseEventFilter):
|
class RoomEncryptedFilter(BaseEventFilter):
|
||||||
"""
|
"""
|
||||||
@@ -10,8 +9,11 @@ class RoomEncryptedFilter(BaseEventFilter):
|
|||||||
def __init__(self):
|
def __init__(self):
|
||||||
super().__init__()
|
super().__init__()
|
||||||
|
|
||||||
async def __call__(self, room: MatrixRoom, event: Event, client: AsyncClient) -> bool:
|
async def __call__(self, context: EventContext) -> bool:
|
||||||
|
if not await super().__call__(context):
|
||||||
|
return False
|
||||||
try:
|
try:
|
||||||
return room.encrypted
|
context[CTX_ROOM_ENCRYPTED] = context.room.encrypted
|
||||||
|
return context.room.encrypted
|
||||||
except:
|
except:
|
||||||
return False
|
return False
|
||||||
@@ -4,15 +4,8 @@ from pathlib import Path
|
|||||||
from dataclasses import dataclass
|
from dataclasses import dataclass
|
||||||
from enum import Enum
|
from enum import Enum
|
||||||
|
|
||||||
from nio import MatrixRoom, Event
|
|
||||||
from nio import UploadResponse
|
from nio import UploadResponse
|
||||||
|
|
||||||
from .filters.base import BaseEventFilter
|
|
||||||
|
|
||||||
from typing import TYPE_CHECKING
|
|
||||||
if TYPE_CHECKING:
|
|
||||||
from .bot import MatrixBot
|
|
||||||
|
|
||||||
@dataclass
|
@dataclass
|
||||||
class MatrixBotConfig:
|
class MatrixBotConfig:
|
||||||
"""Configuration for MatrixBot"""
|
"""Configuration for MatrixBot"""
|
||||||
@@ -68,22 +61,6 @@ class VideoFileProperties:
|
|||||||
thumbnail: Path | str | bytes | None = None
|
thumbnail: Path | str | bytes | None = None
|
||||||
"""Path to the thumbnail or the raw JPEG thumbnail data"""
|
"""Path to the thumbnail or the raw JPEG thumbnail data"""
|
||||||
|
|
||||||
@dataclass
|
|
||||||
class RoomEventData:
|
|
||||||
"""Dataclass that hold information about event that happened in the room"""
|
|
||||||
|
|
||||||
room: MatrixRoom
|
|
||||||
"""The room the event has happened in"""
|
|
||||||
|
|
||||||
event: Event
|
|
||||||
"""The event that has happened in the room"""
|
|
||||||
|
|
||||||
filter: BaseEventFilter
|
|
||||||
"""The filter that invoked this event"""
|
|
||||||
|
|
||||||
bot: "MatrixBot"
|
|
||||||
"""The bot that is the source of the event"""
|
|
||||||
|
|
||||||
@dataclass
|
@dataclass
|
||||||
class UploadResult:
|
class UploadResult:
|
||||||
"""Result of data upload"""
|
"""Result of data upload"""
|
||||||
@@ -109,3 +86,20 @@ class MessageType(Enum):
|
|||||||
AUDIO = "m.audio"
|
AUDIO = "m.audio"
|
||||||
LOCATION = "m.location"
|
LOCATION = "m.location"
|
||||||
VIDEO = "m.video"
|
VIDEO = "m.video"
|
||||||
|
|
||||||
|
class ContextDataKey[T]:
|
||||||
|
"""
|
||||||
|
Instances of this class represent a single possible data key that can be
|
||||||
|
stored inside EventContext.
|
||||||
|
"""
|
||||||
|
def __init__(self, name: str) -> None:
|
||||||
|
"""
|
||||||
|
Initialize a ContextDataKey
|
||||||
|
|
||||||
|
Args:
|
||||||
|
- name - name that will be used internally
|
||||||
|
"""
|
||||||
|
self._name = name
|
||||||
|
|
||||||
|
def __repr__(self) -> str:
|
||||||
|
return f"ContextDataKey[{type(T)}]({repr(self._name)})"
|
||||||
Reference in New Issue
Block a user