init, .gitignore

This commit is contained in:
2025-11-11 00:02:49 +01:00
commit 4e18cd0eb3
124 changed files with 16300 additions and 0 deletions
View File
+32
View File
@@ -0,0 +1,32 @@
import os
import environs
try:
env = environs.Env()
env.read_env("./.env")
except FileNotFoundError:
print("No .env file found, using os.environ.")
api_id = int(os.getenv("API_ID", env.int("API_ID")))
api_hash = os.getenv("API_HASH", env.str("API_HASH"))
STRINGSESSION = os.getenv("STRINGSESSION", env.str("STRINGSESSION"))
second_session = os.getenv("SECOND_SESSION", env.str("SECOND_SESSION", ""))
db_type = os.getenv("DATABASE_TYPE", env.str("DATABASE_TYPE"))
db_url = os.getenv("DATABASE_URL", env.str("DATABASE_URL", ""))
db_name = os.getenv("DATABASE_NAME", env.str("DATABASE_NAME"))
apiflash_key = os.getenv("APIFLASH_KEY", env.str("APIFLASH_KEY"))
rmbg_key = os.getenv("RMBG_KEY", env.str("RMBG_KEY", ""))
vt_key = os.getenv("VT_KEY", env.str("VT_KEY", ""))
gemini_key = os.getenv("GEMINI_KEY", env.str("GEMINI_KEY", ""))
cohere_key = os.getenv("COHERE_KEY", env.str("COHERE_KEY", ""))
pm_limit = int(os.getenv("PM_LIMIT", env.int("PM_LIMIT", 4)))
test_server = bool(os.getenv("TEST_SERVER", env.bool("TEST_SERVER", False)))
modules_repo_branch = os.getenv(
"MODULES_REPO_BRANCH", env.str("MODULES_REPO_BRANCH", "master")
)
+189
View File
@@ -0,0 +1,189 @@
# Moon-Userbot - telegram userbot
# Copyright (C) 2020-present Moon Userbot Organization
#
# This program is free software: you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation, either version 3 of the License, or
# (at your option) any later version.
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU General Public License for more details.
# You should have received a copy of the GNU General Public License
# along with this program. If not, see <https://www.gnu.org/licenses/>.
from collections import OrderedDict
from pyrogram import Client, filters, types
from pyrogram.handlers import MessageHandler
from pyrogram.enums.parse_mode import ParseMode
import asyncio
from typing import Union, List, Dict, Optional
class _TrueFilter(filters.Filter):
async def __call__(self, client: Client, update: types.Message):
return True
class Conversation:
_locks: Dict[int, asyncio.Lock] = {}
def __init__(
self,
client: Client,
chat: Union[str, int],
timeout: float = 5,
delete_at_end=True,
exclusive=True,
):
self.client = client
self.chat = chat
self.timeout = timeout
self.delete_at_end = delete_at_end
self.exclusive = exclusive
self._chat_id = 0
self._message_ids = []
self._handler_object = None
self._chat_unique_lock: Optional[asyncio.Lock] = None
self._waiters: Dict[asyncio.Event, filters.Filter] = {}
self._responses: Dict[asyncio.Event, types.Message] = {}
self._pending_updates: List[types.Message] = []
async def __aenter__(self):
self._chat_id = (await self.client.get_chat(self.chat)).id
if self._chat_id in self._locks:
self._chat_unique_lock = self._locks[self._chat_id]
else:
self._chat_unique_lock = self._locks[self._chat_id] = asyncio.Lock()
if self.exclusive:
await self._chat_unique_lock.acquire()
self._handler_object = MessageHandler(
self._handler, filters.chat(self._chat_id)
)
if -999 not in self.client.dispatcher.groups:
new_groups = OrderedDict(self.client.dispatcher.groups)
new_groups[-999] = []
self.client.dispatcher.groups = new_groups
self.client.dispatcher.groups[-999].append(self._handler_object)
await asyncio.sleep(0)
return self
async def __aexit__(self, exc_type, exc_val, exc_tb):
self.client.dispatcher.groups[-999].remove(self._handler_object)
if self.delete_at_end:
await self.client.delete_messages(self._chat_id, self._message_ids)
if self.exclusive:
self._chat_unique_lock.release()
async def _handler(self, _, message: types.Message):
for event, message_filter in self._waiters.items():
if await message_filter(self.client, message):
self._responses[event] = message
event.set()
break
else:
self._pending_updates.append(message)
message.continue_propagation()
async def get_response(
self,
message_filter: Optional[filters.Filter] = None,
timeout: float = None,
) -> types.Message:
if timeout is None:
timeout = self.timeout
if message_filter is None:
message_filter = _TrueFilter()
for message in self._pending_updates:
if await message_filter(self.client, message):
self._pending_updates.remove(message)
break
else:
message = await self._wait_message(message_filter, timeout)
self._message_ids.append(message.id)
return message
async def _wait_message(
self, message_filter: Optional[filters.Filter], timeout: float
) -> types.Message:
event = asyncio.Event()
self._waiters[event] = message_filter
try:
await asyncio.wait_for(event.wait(), timeout=timeout)
except asyncio.TimeoutError as e:
raise TimeoutError from e
finally:
self._waiters.pop(event)
return self._responses.pop(event)
async def send_message(
self,
text: str,
parse_mode: Optional[str] = ParseMode.HTML,
entities: List[types.MessageEntity] = None,
disable_web_page_preview: bool = None,
disable_notification: bool = None,
reply_to_message_id: int = None,
schedule_date: int = None,
) -> types.Message:
"""Send text messages.
Parameters:
text (``str``):
Text of the message to be sent.
parse_mode (``str``, *optional*):
By default, texts are parsed using HTML style.
Pass "markdown" or "md" to enable Markdown-style parsing.
Pass None to completely disable style parsing.
entities (List of :obj:`~pyrogram.types.MessageEntity`):
List of special entities that appear in message text, which can be specified instead of *parse_mode*.
disable_web_page_preview (``bool``, *optional*):
Disables link previews for links in this message.
disable_notification (``bool``, *optional*):
Sends the message silently.
Users will receive a notification with no sound.
reply_to_message_id (``int``, *optional*):
If the message is a reply, ID of the original message.
schedule_date (``int``, *optional*):
Date when the message will be automatically sent. Unix time.
Returns:
:obj:`~pyrogram.types.Message`: On success, the sent text message is returned.
"""
sent = await self.client.send_message(
chat_id=self._chat_id,
text=text,
parse_mode=parse_mode,
entities=entities,
disable_web_page_preview=disable_web_page_preview,
disable_notification=disable_notification,
reply_to_message_id=reply_to_message_id,
schedule_date=schedule_date,
)
self._message_ids.append(sent.id)
return sent
+237
View File
@@ -0,0 +1,237 @@
# Moon-Userbot - telegram userbot
# Copyright (C) 2020-present Moon Userbot Organization
#
# This program is free software: you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation, either version 3 of the License, or
# (at your option) any later version.
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU General Public License for more details.
# You should have received a copy of the GNU General Public License
# along with this program. If not, see <https://www.gnu.org/licenses/>.
import re
import json
import threading
import sqlite3
from dns import resolver
import pymongo
from utils import config
resolver.default_resolver = resolver.Resolver(configure=False)
resolver.default_resolver.nameservers = ["1.1.1.1"]
class Database:
def get(self, module: str, variable: str, default=None):
"""Get value from database"""
raise NotImplementedError
def set(self, module: str, variable: str, value):
"""Set key in database"""
raise NotImplementedError
def remove(self, module: str, variable: str):
"""Remove key from database"""
raise NotImplementedError
def get_collection(self, module: str) -> dict:
"""Get database for selected module"""
raise NotImplementedError
def close(self):
"""Close the database"""
raise NotImplementedError
class MongoDatabase(Database):
def __init__(self, url, name):
self._client = pymongo.MongoClient(url)
self._database = self._client[name]
def set(self, module: str, variable: str, value):
if not isinstance(module, str) or not isinstance(variable, str):
raise ValueError("Module and variable must be strings")
self._database[module].replace_one(
{"var": variable}, {"var": variable, "val": value}, upsert=True
)
def get(self, module: str, variable: str, default=None):
if not isinstance(module, str) or not isinstance(variable, str):
raise ValueError("Module and variable must be strings")
doc = self._database[module].find_one({"var": variable})
return default if doc is None else doc["val"]
def get_collection(self, module: str):
if not isinstance(module, str):
raise ValueError("Module must be a string")
return {item["var"]: item["val"] for item in self._database[module].find()}
def remove(self, module: str, variable: str):
if not isinstance(module, str) or not isinstance(variable, str):
raise ValueError("Module and variable must be strings")
self._database[module].delete_one({"var": variable})
def close(self):
self._client.close()
def add_chat_history(self, user_id, message):
chat_history = self.get_chat_history(user_id, default=[])
chat_history.append(message)
self.set(f"core.cohere.user_{user_id}", "chat_history", chat_history)
def get_chat_history(self, user_id, default=None):
if default is None:
default = []
return self.get(f"core.cohere.user_{user_id}", "chat_history", default=[])
def addaiuser(self, user_id):
chatai_users = self.get("core.chatbot", "chatai_users", default=[])
if user_id not in chatai_users:
chatai_users.append(user_id)
self.set("core.chatbot", "chatai_users", chatai_users)
def remaiuser(self, user_id):
chatai_users = self.get("core.chatbot", "chatai_users", default=[])
if user_id in chatai_users:
chatai_users.remove(user_id)
self.set("core.chatbot", "chatai_users", chatai_users)
def getaiusers(self):
return self.get("core.chatbot", "chatai_users", default=[])
class SqliteDatabase(Database):
def __init__(self, file):
self._conn = sqlite3.connect(file, check_same_thread=False)
self._conn.row_factory = sqlite3.Row
self._cursor = self._conn.cursor()
self._lock = threading.Lock()
@staticmethod
def _parse_row(row: sqlite3.Row):
if row["type"] == "bool":
return row["val"] == "1"
if row["type"] == "int":
return int(row["val"])
if row["type"] == "str":
return row["val"]
return json.loads(row["val"])
def _execute(self, module: str, *args, **kwargs) -> sqlite3.Cursor:
pattern = r"^(core|custom)"
if not re.match(pattern, module):
raise ValueError(f"Invalid module name format: {module}")
self._lock.acquire()
try:
cursor = self._conn.cursor()
return cursor.execute(*args, **kwargs)
except sqlite3.OperationalError as e:
if str(e).startswith("no such table"):
sql = f"""
CREATE TABLE IF NOT EXISTS '{module}' (
var TEXT UNIQUE NOT NULL,
val TEXT NOT NULL,
type TEXT NOT NULL
)
"""
cursor = self._conn.cursor()
cursor.execute(sql)
self._conn.commit()
return cursor.execute(*args, **kwargs)
raise e from None
finally:
self._lock.release()
def get(self, module: str, variable: str, default=None):
sql = f"SELECT * FROM '{module}' WHERE var=?"
cur = self._execute(module, sql, (variable,))
row = cur.fetchone()
if row is None:
return default
return self._parse_row(row)
def set(self, module: str, variable: str, value) -> bool:
sql = f"""
INSERT INTO '{module}' VALUES ( ?, ?, ? )
ON CONFLICT (var) DO
UPDATE SET val=?, type=? WHERE var=?
"""
if isinstance(value, bool):
val = "1" if value else "0"
typ = "bool"
elif isinstance(value, str):
val = value
typ = "str"
elif isinstance(value, int):
val = str(value)
typ = "int"
else:
val = json.dumps(value)
typ = "json"
self._execute(module, sql, (variable, val, typ, val, typ, variable))
self._conn.commit()
return True
def remove(self, module: str, variable: str):
sql = f"DELETE FROM '{module}' WHERE var=?"
self._execute(module, sql, (variable,))
self._conn.commit()
def get_collection(self, module: str) -> dict:
pattern = r"^(core|custom)"
if not re.match(pattern, module):
raise ValueError(f"Invalid module name format: {module}")
sql = f"SELECT * FROM '{module}'"
cur = self._execute(module, sql)
collection = {}
for row in cur:
collection[row["var"]] = self._parse_row(row)
return collection
def close(self):
self._conn.commit()
self._conn.close()
def add_chat_history(self, user_id, message):
chat_history = self.get_chat_history(user_id, default=[])
chat_history.append(message)
self.set(f"core.cohere.user_{user_id}", "chat_history", chat_history)
def get_chat_history(self, user_id, default=None):
if default is None:
default = []
return self.get(f"core.cohere.user_{user_id}", "chat_history", default=[])
def addaiuser(self, user_id):
chatai_users = self.get("core.chatbot", "chatai_users", default=[])
if user_id not in chatai_users:
chatai_users.append(user_id)
self.set("core.chatbot", "chatai_users", chatai_users)
def remaiuser(self, user_id):
chatai_users = self.get("core.chatbot", "chatai_users", default=[])
if user_id in chatai_users:
chatai_users.remove(user_id)
self.set("core.chatbot", "chatai_users", chatai_users)
def getaiusers(self):
return self.get("core.chatbot", "chatai_users", default=[])
if config.db_type in ["mongo", "mongodb"]:
db = MongoDatabase(config.db_url, config.db_name)
else:
db = SqliteDatabase(config.db_name)
+1251
View File
File diff suppressed because it is too large Load Diff
+55
View File
@@ -0,0 +1,55 @@
# Moon-Userbot - telegram userbot
# Copyright (C) 2020-present Moon Userbot Organization
#
# This program is free software: you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation, either version 3 of the License, or
# (at your option) any later version.
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU General Public License for more details.
# You should have received a copy of the GNU General Public License
# along with this program. If not, see <https://www.gnu.org/licenses/>.
from sys import version_info
from .db import db
import git
__all__ = [
"modules_help",
"requirements_list",
"python_version",
"prefix",
"gitrepo",
"userbot_version",
]
modules_help = {}
requirements_list = []
python_version = f"{version_info[0]}.{version_info[1]}.{version_info[2]}"
prefix = db.get("core.main", "prefix", ".")
try:
gitrepo = git.Repo(".")
except git.exc.InvalidGitRepositoryError:
repo = git.Repo.init()
origin = repo.create_remote(
"origin", "https://github.com/The-MoonTg-project/Moon-Userbot"
)
origin.fetch()
repo.create_head("main", origin.refs.main)
repo.heads.main.set_tracking_branch(origin.refs.main)
repo.heads.main.checkout(True)
gitrepo = git.Repo(".")
if len(gitrepo.tags) > 0:
commits_since_tag = list(gitrepo.iter_commits(f"{gitrepo.tags[-1].name}..HEAD"))
else:
commits_since_tag = []
userbot_version = f"2.5.{len(commits_since_tag)}"
+78
View File
@@ -0,0 +1,78 @@
import logging
from pathlib import Path
from typing import Optional
from pyrogram import Client
from utils.scripts import load_module
from utils.misc import modules_help
class ModuleManager:
_instance: Optional["ModuleManager"] = None
def __init__(self):
self.success_modules = 0
self.failed_modules = 0
self.help_navigator = None
@classmethod
def get_instance(cls) -> "ModuleManager":
if cls._instance is None:
cls._instance = ModuleManager()
return cls._instance
async def load_modules(self, app: Client):
"""Load all modules and initialize help navigator"""
for path in Path("modules").rglob("*.py"):
try:
await load_module(
path.stem, app, core="custom_modules" not in path.parent.parts
)
except Exception:
logging.warning("Can't import module %s", path.stem, exc_info=True)
self.failed_modules += 1
else:
self.success_modules += 1
logging.info("Imported %d modules", self.success_modules)
if self.failed_modules:
logging.warning("Failed to import %d modules", self.failed_modules)
self.help_navigator = HelpNavigator()
return self.help_navigator
class HelpNavigator:
def __init__(self):
self.current_page = 1
self.module_list = list(modules_help.keys())
self.total_pages = (len(modules_help) + 9) // 10
logging.info("Initialized HelpNavigator with %d modules", len(self.module_list))
async def send_page(self, message):
from utils.misc import prefix
start_index = (self.current_page - 1) * 10
end_index = start_index + 10
page_modules = self.module_list[start_index:end_index]
text = "<b>Help</b>\n"
text += f"For more help on how to use a command, type <code>{prefix}help [module]</code>\n\n"
text += f"Help Page No: {self.current_page}/{self.total_pages}\n\n"
for module_name in page_modules:
commands = modules_help[module_name]
text += f"<b>• {module_name.title()}:</b> {', '.join([f'<code>{prefix + cmd_name.split()[0]}</code>' for cmd_name in commands.keys()])}\n"
text += f"\n<b>The number of modules in the userbot: {len(modules_help)}</b>"
await message.edit(text, disable_web_page_preview=True)
def next_page(self) -> bool:
if self.current_page < self.total_pages:
self.current_page += 1
return True
return False
def prev_page(self) -> bool:
if self.current_page > 1:
self.current_page -= 1
return True
return False
+200
View File
@@ -0,0 +1,200 @@
#!/usr/bin/env python3
# @source: https://github.com/radude/rentry/blob/master/rentry.py
import asyncio
import http.cookiejar
import urllib.parse
import urllib.request
from uuid import uuid4
from datetime import datetime, timedelta
from http.cookies import SimpleCookie
from json import loads as json_loads
from utils.db import db
BASE_PROTOCOL = "https://"
BASE_URL = "rentry.co"
_headers = {"Referer": f"{BASE_PROTOCOL}{BASE_URL}"}
class UrllibClient:
"""Simple HTTP Session Client, keeps cookies."""
def __init__(self):
self.cookie_jar = http.cookiejar.CookieJar()
self.opener = urllib.request.build_opener(
urllib.request.HTTPCookieProcessor(self.cookie_jar)
)
urllib.request.install_opener(self.opener)
def get(self, url, headers=None):
if headers is None:
headers = {}
request = urllib.request.Request(url, headers=headers)
return self._request(request)
def post(self, url, data=None, headers=None):
if headers is None:
headers = {}
postdata = urllib.parse.urlencode(data).encode()
request = urllib.request.Request(url, postdata, headers)
return self._request(request)
def _request(self, request):
response = self.opener.open(request)
response.status_code = response.getcode()
response.data = response.read().decode("utf-8")
return response
def raw(url: str):
client = UrllibClient()
return json_loads(client.get(f"{BASE_PROTOCOL}{BASE_URL}/api/raw/{url}").data)
def new(text: str, edit_code: str = "", url: str = ""):
client, cookie = UrllibClient(), SimpleCookie()
cookie.load(vars(client.get(f"{BASE_PROTOCOL}{BASE_URL}"))["headers"]["Set-Cookie"])
csrftoken = cookie["csrftoken"].value
payload = {
"csrfmiddlewaretoken": csrftoken,
"url": url,
"edit_code": edit_code,
"text": text,
}
return json_loads(
client.post(
f"{BASE_PROTOCOL}{BASE_URL}" + "/api/new", payload, headers=_headers
).data
)
def edit(url_short: str, edit_code: str, text: str):
client, cookie = UrllibClient(), SimpleCookie()
cookie.load(vars(client.get(f"{BASE_PROTOCOL}{BASE_URL}"))["headers"]["Set-Cookie"])
csrftoken = cookie["csrftoken"].value
payload = {"csrfmiddlewaretoken": csrftoken, "edit_code": edit_code, "text": text}
return json_loads(
client.post(
f"{BASE_PROTOCOL}{BASE_URL}/api/edit/{url_short}", payload, headers=_headers
).data
)
def delete(url_short: str, edit_code: str):
client, cookie = UrllibClient(), SimpleCookie()
cookie.load(vars(client.get(f"{BASE_PROTOCOL}{BASE_URL}"))["headers"]["Set-Cookie"])
csrftoken = cookie["csrftoken"].value
payload = {"csrfmiddlewaretoken": csrftoken, "edit_code": edit_code}
return json_loads(
client.post(
f"{BASE_PROTOCOL}{BASE_URL}/api/delete/{url_short}",
payload,
headers=_headers,
).data
)
async def paste(
text: str,
return_edit: bool = False,
edit_bin: bool = False,
edit_code: str = None,
url: str = None,
permanent: bool = False,
) -> str | tuple[str, str]:
"""Pastes some text to rentry bin.
args:
text: Input text to paste
return_edit: If it should return edit code also
edit_bin: If this request is to edit an already existing bin
edit_code: Only required if edit_bin is True. It is the edit code used to edit bin.
url: Only required if edit_bin is True. It is the url on which the bin is located.
permanent: If the pasted content should not be deleted automatically
returns:
The url of the paste or return a tuple containing url and edit code
"""
if not str(text):
return
if edit_bin:
if not (url and edit_code):
raise ValueError("Please provide both, url and edit code")
response = edit(url_short=url, edit_code=edit_code, text=text)
else:
response = new(text=text)
if response.get("status") != "200":
raise RuntimeError(
f"paste task terminated with status: {response.get('status')}\n"
f"Message: {response.get('content', 'No message provided')}"
)
url = response["url"]
edit_code = response["edit_code"]
if not permanent:
short_url = response["url_short"]
time_now = datetime.now()
ftime = time_now.strftime("%d %I:%M:%S %p %Y")
print(f"URL: {url} - Edit Code: {edit_code} - Time: {ftime}")
rallUrls = db.get("core.rentry", "urls", default={"allUrls": {}})
entry_id = str(uuid4())
rallUrls["allUrls"][entry_id] = {
"url": short_url,
"edit_code": edit_code,
"time": ftime,
}
db.set("core.rentry", "urls", rallUrls)
if return_edit:
return (url, edit_code)
return url
async def rentry_cleanup_job():
"""Periodically checks and deletes rentry pastes older than 24 hours"""
while True:
try:
rallUrls = db.get("core.rentry", "urls", default={"allUrls": {}})
now = datetime.now()
deleted_count = 0
error_count = 0
for entry_id, entry in list(rallUrls["allUrls"].items()):
url = entry["url"]
entry_time = datetime.strptime(entry["time"], "%d %I:%M:%S %p %Y")
if now - entry_time > timedelta(days=1):
try:
delete(url, entry["edit_code"])
del rallUrls["allUrls"][entry_id]
deleted_count += 1
print(f"[#] Deleted expired rentry paste: {url}")
except Exception as e:
error_count += 1
print(f"[!] Failed to delete rentry paste {url}: {str(e)}")
if deleted_count or error_count:
print(
f"[*] Cleanup summary: {deleted_count} deleted, {error_count} failed"
)
if deleted_count:
db.set("core.rentry", "urls", rallUrls)
except Exception as e:
print(f"[!] Error in rentry cleanup job: {str(e)}")
await asyncio.sleep(12 * 60 * 60)
+542
View File
@@ -0,0 +1,542 @@
# Moon-Userbot - telegram userbot
# Copyright (C) 2020-present Moon Userbot Organization
#
# This program is free software: you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation, either version 3 of the License, or
# (at your option) any later version.
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU General Public License for more details.
# You should have received a copy of the GNU General Public License
# along with this program. If not, see <https://www.gnu.org/licenses/>.
import asyncio
import importlib
import math
import os
import re
import shlex
import subprocess
import sys
import time
import traceback
from PIL import Image
from io import BytesIO
from types import ModuleType
from typing import Dict, Tuple
import psutil
from pyrogram import Client, errors, filters
from pyrogram.errors import FloodWait, MessageNotModified, UserNotParticipant
from pyrogram.types import Message
from pyrogram.enums import ChatMembersFilter
from utils.db import db
from .misc import modules_help, prefix, requirements_list
META_COMMENTS = re.compile(r"^ *# *meta +(\S+) *: *(.*?)\s*$", re.MULTILINE)
interact_with_to_delete = []
def time_formatter(milliseconds: int) -> str:
"""Time Formatter"""
seconds, milliseconds = divmod(int(milliseconds), 1000)
minutes, seconds = divmod(seconds, 60)
hours, minutes = divmod(minutes, 60)
days, hours = divmod(hours, 24)
tmp = (
((str(days) + " day(s), ") if days else "")
+ ((str(hours) + " hour(s), ") if hours else "")
+ ((str(minutes) + " minute(s), ") if minutes else "")
+ ((str(seconds) + " second(s), ") if seconds else "")
+ ((str(milliseconds) + " millisecond(s), ") if milliseconds else "")
)
return tmp[:-2]
def humanbytes(size):
"""Convert Bytes To Bytes So That Human Can Read It"""
if not size:
return ""
power = 2**10
raised_to_pow = 0
dict_power_n = {0: "", 1: "Ki", 2: "Mi", 3: "Gi", 4: "Ti"}
while size > power:
size /= power
raised_to_pow += 1
return str(round(size, 2)) + " " + dict_power_n[raised_to_pow] + "B"
async def edit_or_send_as_file(
tex: str,
message: Message,
client: Client,
caption: str = "<code>Result!</code>",
file_name: str = "result",
):
"""Send As File If Len Of Text Exceeds Tg Limit Else Edit Message"""
if not tex:
await message.edit("<code>Wait, What?</code>")
return
if len(tex) > 1024:
await message.edit("<code>OutPut is Too Large, Sending As File!</code>")
file_names = f"{file_name}.txt"
with open(file_names, "w") as fn:
fn.write(tex)
await client.send_document(message.chat.id, file_names, caption=caption)
await message.delete()
if os.path.exists(file_names):
os.remove(file_names)
return
return await message.edit(tex)
def get_text(message: Message) -> None | str:
"""Extract Text From Commands"""
text_to_return = message.text
if message.text is None:
return None
if " " in text_to_return:
try:
return message.text.split(None, 1)[1]
except IndexError:
return None
else:
return None
async def progress(current, total, message, start, type_of_ps, file_name=None):
"""Progress Bar For Showing Progress While Uploading / Downloading File - Normal"""
now = time.time()
diff = now - start
if round(diff % 10.00) == 0 or current == total:
percentage = current * 100 / total
speed = current / diff
elapsed_time = round(diff) * 1000
if elapsed_time == 0:
return
time_to_completion = round((total - current) / speed) * 1000
estimated_total_time = elapsed_time + time_to_completion
progress_str = f"{''.join(['' for i in range(math.floor(percentage / 10))])}"
progress_str += (
f"{''.join(['' for i in range(10 - math.floor(percentage / 10))])}"
)
progress_str += f"{round(percentage, 2)}%\n"
tmp = f"{progress_str}{humanbytes(current)} of {humanbytes(total)}\n"
tmp += f"ETA: {time_formatter(estimated_total_time)}"
if file_name:
try:
await message.edit(
f"{type_of_ps}\n<b>File Name:</b> <code>{file_name}</code>\n{tmp}"
)
except FloodWait as e:
await asyncio.sleep(e.x)
except MessageNotModified:
pass
else:
try:
await message.edit(f"{type_of_ps}\n{tmp}")
except FloodWait as e:
await asyncio.sleep(e.x)
except MessageNotModified:
pass
async def run_cmd(prefix: str) -> Tuple[str, str, int, int]:
"""Run Commands"""
args = shlex.split(prefix)
process = await asyncio.create_subprocess_exec(
*args, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE
)
stdout, stderr = await process.communicate()
return (
stdout.decode("utf-8", "replace").strip(),
stderr.decode("utf-8", "replace").strip(),
process.returncode,
process.pid,
)
def mediainfo(media):
xx = str((str(media)).split("(", maxsplit=1)[0])
m = ""
if xx == "MessageMediaDocument":
mim = media.document.mime_type
if mim == "application/x-tgsticker":
m = "sticker animated"
elif "image" in mim:
if mim == "image/webp":
m = "sticker"
elif mim == "image/gif":
m = "gif as doc"
else:
m = "pic as doc"
elif "video" in mim:
if "DocumentAttributeAnimated" in str(media):
m = "gif"
elif "DocumentAttributeVideo" in str(media):
i = str(media.document.attributes[0])
if "supports_streaming=True" in i:
m = "video"
m = "video as doc"
else:
m = "video"
elif "audio" in mim:
m = "audio"
else:
m = "document"
elif xx == "MessageMediaPhoto":
m = "pic"
elif xx == "MessageMediaWebPage":
m = "web"
return m
async def edit_or_reply(message, txt):
"""Edit Message If It's From Self, Else Reply To Message"""
if not message:
return
if message.from_user and message.from_user.is_self:
return await message.edit(txt)
return await message.reply(txt)
def text(message: Message) -> str:
"""Find text in `Message` object"""
return message.text if message.text else message.caption
def restart() -> None:
music_bot_pid = db.get("custom.musicbot", "music_bot_pid", None)
if music_bot_pid is not None:
try:
music_bot_process = psutil.Process(music_bot_pid)
music_bot_process.terminate()
except psutil.NoSuchProcess:
print("Music bot is not running.")
os.execvp(sys.executable, [sys.executable, "main.py"]) # skipcq
def format_exc(e: Exception, suffix="") -> str:
traceback.print_exc()
err = traceback.format_exc()
if isinstance(e, errors.RPCError):
return (
f"<b>Telegram API error!</b>\n"
f"<code>[{e.CODE} {e.ID or e.NAME}] — {e.MESSAGE.format(value=e.value)}</code>\n\n<b>{suffix}</b>"
)
return f"<b>Error!</b>\n" f"<code>{err}</code>"
def with_reply(func):
async def wrapped(client: Client, message: Message):
if not message.reply_to_message:
await message.edit("<b>Reply to message is required</b>")
else:
return await func(client, message)
return wrapped
def is_admin(func):
async def wrapped(client: Client, message: Message):
try:
chat_member = await client.get_chat_member(
message.chat.id, message.from_user.id
)
chat_admins = []
async for member in client.get_chat_members(
message.chat.id, filter=ChatMembersFilter.ADMINISTRATORS
):
chat_admins.append(member)
if chat_member in chat_admins:
return await func(client, message)
await message.edit("You need to be an admin to perform this action.")
except UserNotParticipant:
await message.edit(
"You need to be a participant in the chat to perform this action."
)
return wrapped
async def interact_with(message: Message) -> Message:
"""
Check history with bot and return bot's response
Example:
.. code-block:: python
bot_msg = await interact_with(await bot.send_message("@BotFather", "/start"))
:param message: already sent message to bot
:return: bot's response
"""
await asyncio.sleep(1)
# noinspection PyProtectedMember
response = [
msg async for msg in message._client.get_chat_history(message.chat.id, limit=1)
]
seconds_waiting = 0
while response[0].from_user.is_self:
seconds_waiting += 1
if seconds_waiting >= 5:
raise RuntimeError("bot didn't answer in 5 seconds")
await asyncio.sleep(1)
# noinspection PyProtectedMember
response = [
msg
async for msg in message._client.get_chat_history(message.chat.id, limit=1)
]
interact_with_to_delete.append(message.id)
interact_with_to_delete.append(response[0].id)
return response[0]
def format_module_help(module_name: str, full=True):
commands = modules_help[module_name]
help_text = (
f"<b>Help for |{module_name}|\n\nUsage:</b>\n" if full else "<b>Usage:</b>\n"
)
for command, desc in commands.items():
cmd = command.split(maxsplit=1)
args = " <code>" + cmd[1] + "</code>" if len(cmd) > 1 else ""
help_text += f"<code>{prefix}{cmd[0]}</code>{args} — <i>{desc}</i>\n"
return help_text
def format_small_module_help(module_name: str, full=True):
commands = modules_help[module_name]
help_text = (
f"<b>Help for |{module_name}|\n\nCommands list:\n"
if full
else "<b>Commands list:\n"
)
for command, _desc in commands.items():
cmd = command.split(maxsplit=1)
args = " <code>" + cmd[1] + "</code>" if len(cmd) > 1 else ""
help_text += f"<code>{prefix}{cmd[0]}</code>{args}\n"
help_text += f"\nGet full usage: <code>{prefix}help {module_name}</code></b>"
return help_text
def import_library(library_name: str, package_name: str = None):
"""
Loads a library, or installs it in ImportError case
:param library_name: library name (import example...)
:param package_name: package name in PyPi (pip install example)
:return: loaded module
"""
if package_name is None:
package_name = library_name
requirements_list.append(package_name)
try:
return importlib.import_module(library_name)
except ImportError as exc:
completed = subprocess.run(
[sys.executable, "-m", "pip", "install", "--upgrade", package_name],
check=True,
)
if completed.returncode != 0:
raise AssertionError(
f"Failed to install library {package_name} (pip exited with code {completed.returncode})"
) from exc
return importlib.import_module(library_name)
def uninstall_library(package_name: str):
"""
Uninstalls a library
:param package_name: package name in PyPi (pip uninstall example)
"""
completed = subprocess.run(
[sys.executable, "-m", "pip", "uninstall", "-y", package_name], check=True
)
if completed.returncode != 0:
raise AssertionError(
f"Failed to uninstall library {package_name} (pip exited with code {completed.returncode})"
)
def resize_image(
input_img, output=None, img_type="PNG", size: int = 512, size2: int = None
):
if output is None:
output = BytesIO()
output.name = f"sticker.{img_type.lower()}"
with Image.open(input_img) as img:
# We used to use thumbnail(size) here, but it returns with a *max* dimension of 512,512
# rather than making one side exactly 512, so we have to calculate dimensions manually :(
if size2 is not None:
size = (size, size2)
elif img.width == img.height:
size = (size, size)
elif img.width < img.height:
size = (max(size * img.width // img.height, 1), size)
else:
size = (size, max(size * img.height // img.width, 1))
img.resize(size).save(output, img_type)
return output
def resize_new_image(image_path, output_path, desired_width=None, desired_height=None):
"""
Resize an image to the desired dimensions while maintaining the aspect ratio.
Args:
image_path (str): Path to the input image file.
output_path (str): Path to save the resized image.
desired_width (int, optional): Desired width in pixels. If not provided, the aspect ratio will be maintained.
desired_height (int, optional): Desired height in pixels. If not provided, the aspect ratio will be maintained.
"""
image = Image.open(image_path)
width, height = image.size
aspect_ratio = width / height
if desired_width and desired_height:
new_width, new_height = desired_width, desired_height
elif desired_height:
new_width, new_height = int(desired_height * aspect_ratio), desired_height
else:
new_width, new_height = 150, 150
resized_image = image.resize((new_width, new_height), Image.Resampling.LANCZOS)
resized_image.save(output_path)
if os.path.exists(image_path):
os.remove(image_path)
async def load_module(
module_name: str,
client: Client,
message: Message = None,
core=False,
) -> ModuleType:
if module_name in modules_help and not core:
await unload_module(module_name, client)
path = f"modules.{'custom_modules.' if not core else ''}{module_name}"
with open(f"{path.replace('.', '/')}.py", encoding="utf-8") as f:
code = f.read()
meta = parse_meta_comments(code)
packages = meta.get("requires", "").split()
requirements_list.extend(packages)
try:
module = importlib.import_module(path)
except ImportError as e:
if core:
# Core modules shouldn't raise ImportError
raise
if not packages:
raise
if message:
await message.edit(f"<b>Installing requirements: {' '.join(packages)}</b>")
proc = await asyncio.create_subprocess_exec(
sys.executable,
"-m",
"pip",
"install",
"-U",
*packages,
)
try:
await asyncio.wait_for(proc.wait(), timeout=120)
except asyncio.TimeoutError:
if message:
await message.edit(
"<b>Timeout while installed requirements."
+ "Try to install them manually</b>"
)
raise TimeoutError("timeout while installing requirements") from e
if proc.returncode != 0:
if message:
await message.edit(
f"<b>Failed to install requirements (pip exited with code {proc.returncode}). "
f"Check logs for futher info</b>",
)
raise RuntimeError("failed to install requirements") from e
module = importlib.import_module(path)
for _name, obj in vars(module).items():
if isinstance(getattr(obj, "handlers", []), list):
for handler, group in getattr(obj, "handlers", []):
client.add_handler(handler, group)
module.__meta__ = meta
return module
async def unload_module(module_name: str, client: Client) -> bool:
path = "modules.custom_modules." + module_name
if path not in sys.modules:
return False
module = importlib.import_module(path)
for _name, obj in vars(module).items():
for handler, group in getattr(obj, "handlers", []):
client.remove_handler(handler, group)
del modules_help[module_name]
del sys.modules[path]
return True
def no_prefix(handler):
def func(_, __, message):
if message.text and not message.text.startswith(handler):
return True
return False
return filters.create(func)
def parse_meta_comments(code: str) -> Dict[str, str]:
try:
groups = META_COMMENTS.search(code).groups()
except AttributeError:
return {}
return {groups[i]: groups[i + 1] for i in range(0, len(groups), 2)}
def ReplyCheck(message: Message):
reply_id = None
if message.reply_to_message:
reply_id = message.reply_to_message.id
elif not message.from_user.is_self:
reply_id = message.id
return reply_id