Asynchronous commands
This commit is contained in:
parent
9f5cec705e
commit
32a185060b
6 changed files with 88 additions and 38 deletions
8
main.py
8
main.py
|
|
@ -5,7 +5,8 @@ import asyncio
|
|||
|
||||
from scheduling import OrderScheduler
|
||||
from orders import order_issue, order_check
|
||||
from telegram import handle_commands
|
||||
from telegram.telegram import handle_commands
|
||||
from telegram.commands import commands
|
||||
from db.queries import initdb
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
|
@ -41,5 +42,8 @@ if __name__=='__main__':
|
|||
else:
|
||||
loop = asyncio.new_event_loop()
|
||||
s = OrderScheduler(loop)
|
||||
loop.run_until_complete(handle_commands())
|
||||
loop.run_until_complete(handle_commands(
|
||||
commands=commands,
|
||||
loop=loop
|
||||
))
|
||||
loop.run_forever()
|
||||
|
|
|
|||
|
|
@ -5,7 +5,7 @@ from util import make_session
|
|||
from generate import generate_order, generate_punishment
|
||||
from db.queries import order_status_put, punishment_status_put, order_status_outstanding, order_status_confirm
|
||||
from mastodon import Mastodon
|
||||
from telegram import Telegram
|
||||
from telegram.telegram import Telegram
|
||||
from settings import MASTODON_USERNAME, ORDER_TIMEOUT, ENV
|
||||
from util import timezone
|
||||
|
||||
|
|
|
|||
|
|
@ -15,6 +15,7 @@ MASTODON_VISIBILITY = os.environ.get('MASTODON_VISIBILITY', 'direct')
|
|||
|
||||
TELEGRAM_API_TOKEN = os.environ.get('TELEGRAM_API_TOKEN')
|
||||
TELEGRAM_CHAT_ID = int(os.environ.get('TELEGRAM_CHAT_ID'))
|
||||
TELEGRAM_COMMAND_TIMEOUT = int(os.environ.get('TELEGRAM_COMMAND_TIMEOUT', 120))
|
||||
|
||||
SQLITE_DB = os.environ.get('SQLITE_DB', 'db.sqlite3')
|
||||
|
||||
|
|
|
|||
0
telegram/__init__.py
Normal file
0
telegram/__init__.py
Normal file
31
telegram/commands.py
Normal file
31
telegram/commands.py
Normal file
|
|
@ -0,0 +1,31 @@
|
|||
import re
|
||||
import logging
|
||||
import peewee
|
||||
|
||||
from db.queries import skip_day_put
|
||||
from .telegram import TelegramCommand
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
class SkipDayAddCommand(TelegramCommand):
|
||||
command_text = "/skip_day"
|
||||
date_regex = re.compile(r"^(?P<date>\d{4}-\d{2}-\d{2})$")
|
||||
|
||||
async def exec_inner(self, text, update, session, forResponse=None, reply=None):
|
||||
try:
|
||||
yield "Please enter a skip day"
|
||||
|
||||
response = await forResponse()
|
||||
|
||||
m = self.date_regex.match(response)
|
||||
if(m is not None):
|
||||
skip_day_put(m.group('date'))
|
||||
yield f"Skip day {m.group('date')} has been added"
|
||||
else:
|
||||
yield "Please enter a valid date"
|
||||
except peewee.IntegrityError:
|
||||
yield "That day has already been added"
|
||||
|
||||
commands = [
|
||||
SkipDayAddCommand()
|
||||
]
|
||||
|
|
@ -1,9 +1,7 @@
|
|||
import re
|
||||
import logging
|
||||
import asyncio
|
||||
|
||||
from settings import TELEGRAM_API_TOKEN, TELEGRAM_CHAT_ID
|
||||
from db.queries import skip_day_put
|
||||
from settings import TELEGRAM_API_TOKEN, TELEGRAM_CHAT_ID, TELEGRAM_COMMAND_TIMEOUT
|
||||
from util import make_session
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
|
@ -66,33 +64,28 @@ class TelegramCommand():
|
|||
(self.command_regex is not None and self.command_regex.match(text) is not None)
|
||||
)
|
||||
|
||||
async def exec(self, text, update, session):
|
||||
async def exec(self, *args, **kwargs):
|
||||
try:
|
||||
async for reply in self.exec_inner(*args, **kwargs):
|
||||
yield reply
|
||||
except Exception as e:
|
||||
yield "There was a problem with your command"
|
||||
logger.error(f"{self.__class__.__name__} {e}")
|
||||
|
||||
async def exec_inner(self, text, update, session):
|
||||
pass
|
||||
|
||||
class SkipDayAddCommand(TelegramCommand):
|
||||
command_regex = re.compile(r"^\/skip_day( (?P<date>\d{4}-\d{2}-\d{2}))?$")
|
||||
|
||||
async def exec(self, text, update, session):
|
||||
date_str = text.split(' ')[1]
|
||||
|
||||
skip_day_put(date_str)
|
||||
|
||||
yield f"Added skip day {date_str}"
|
||||
|
||||
commands = [
|
||||
SkipDayAddCommand()
|
||||
]
|
||||
|
||||
async def do_reply(data, t):
|
||||
if(type(data) is str):
|
||||
await t.message_send(
|
||||
data
|
||||
)
|
||||
|
||||
async def handle_commands(commands=commands):
|
||||
async def handle_commands(commands=[], loop=None):
|
||||
async with make_session() as session:
|
||||
t = Telegram(session)
|
||||
|
||||
command_tasks = set()
|
||||
command_futures = {}
|
||||
def forResponse():
|
||||
f = loop.create_future()
|
||||
command_futures[chat_id] = f
|
||||
return f
|
||||
|
||||
offset = 0
|
||||
while True:
|
||||
try:
|
||||
|
|
@ -106,20 +99,41 @@ async def handle_commands(commands=commands):
|
|||
if(chat_id != TELEGRAM_CHAT_ID):
|
||||
continue
|
||||
|
||||
if(chat_id in command_futures):
|
||||
command_futures[chat_id].set_result(
|
||||
update['message']['text']
|
||||
)
|
||||
del command_futures[chat_id]
|
||||
continue
|
||||
|
||||
for command in commands:
|
||||
if(command.matches(update['message']['text'])):
|
||||
async def timeoutTask():
|
||||
try:
|
||||
result = command.exec(update['message']['text'], update, session)
|
||||
async with asyncio.timeout(TELEGRAM_COMMAND_TIMEOUT):
|
||||
async for reply in command.exec(
|
||||
update['message']['text'],
|
||||
update,
|
||||
session,
|
||||
forResponse=forResponse
|
||||
):
|
||||
await t.message_send(reply)
|
||||
except TimeoutError:
|
||||
if chat_id in command_futures:
|
||||
del command_futures[chat_id]
|
||||
|
||||
if result.__class__.__name__ == 'async_generator':
|
||||
async for r in result:
|
||||
await do_reply(r, t)
|
||||
else:
|
||||
await do_reply(await result, t)
|
||||
except:
|
||||
await t.message_send("Your command has timed out")
|
||||
except Exception:
|
||||
await t.message_send('There was a problem with your command')
|
||||
logger.exception('Problem while executing a command')
|
||||
|
||||
task = asyncio.create_task(
|
||||
timeoutTask()
|
||||
)
|
||||
|
||||
command_tasks.add(task)
|
||||
task.add_done_callback(command_tasks.discard)
|
||||
|
||||
break
|
||||
offset = update['update_id'] + 1
|
||||
else:
|
||||
Loading…
Reference in a new issue