diff --git a/deploy/aggregator.service b/deploy/aggregator.service index d4d117b..2808ce4 100644 --- a/deploy/aggregator.service +++ b/deploy/aggregator.service @@ -1,6 +1,6 @@ [Unit] Description=Makerspace Leiden Aggregator -After=network.target redis.service signal.service +After=network.target redis.service [Service] Type=simple diff --git a/requirements.txt b/requirements.txt index d0c9fd4..af4eaaf 100644 --- a/requirements.txt +++ b/requirements.txt @@ -5,9 +5,6 @@ aiocron==1.6 # Development tools croniter==1.0.13 -# Signal BOT -dbussy==1.3 - # Human-readable time deltas humanize==3.5.0 @@ -19,8 +16,6 @@ pytest pytest-cov python-daemon==2.3.0 -# Telegram BOT -python-telegram-bot==13.5 quart==0.15.1 redis==3.5.3 requests==2.31.0 diff --git a/server-dev.py b/server-dev.py index aa6121e..b854ba1 100755 --- a/server-dev.py +++ b/server-dev.py @@ -49,11 +49,6 @@ "email": { "from_address": "MakerSpace BOT ", }, - # 'telegram_bot': { - # 'api_token': os.environ['TELEGRAM_BOT_API'], - # }, - # 'signal_bot': { - # }, } diff --git a/server-prod.py b/server-prod.py index 6d58e15..9394a01 100755 --- a/server-prod.py +++ b/server-prod.py @@ -37,7 +37,6 @@ "key_prefix": "msl", "users_expiration_time_in_sec": 60, "pending_machine_activation_timeout_in_sec": 90, - "telegram_token_expiration_in_sec": 5 * 60, # 5 minutes "machine_state_timeout_in_minutes": 60, # 1 hour "history_lines_expiration_in_days": 7, }, @@ -68,12 +67,6 @@ "email": { "from_address": "MakerSpace BOT ", }, - "telegram_bot": { - "api_token": os.environ["TELEGRAM_BOT_API"], - }, - "signal_bot": { - "some": "config", - }, } diff --git a/src/aggregator/bots/bot_logic.py b/src/aggregator/bots/bot_logic.py deleted file mode 100644 index f33d3c9..0000000 --- a/src/aggregator/bots/bot_logic.py +++ /dev/null @@ -1,128 +0,0 @@ -from aggregator.messages import ( - BASIC_COMMANDS, - COMMAND_CHECKIN, - COMMAND_HELP, - COMMAND_NO, - COMMAND_OUT, - COMMAND_WHO, - COMMAND_YES, - STATE_CONFIRM_CHECKOUT, - STATE_CONFIRM_VOLUNTEERING, - MessageCancelAction, - MessageConfirmCheckout, - MessageConfirmedCheckout, - MessageConfirmedVolunteering, - MessageHelp, - MessageNotRegistered, - MessageUnknown, - MessageUserNotInSpace, - MessageVolunteeringNotNecessary, - MessageWho, -) - - -class BotLogic(object): - def __init__(self, aggregator): - self.aggregator = aggregator - self.chat_states = ChatStates(aggregator.clock) - - def handle_new_conversation(self, chat_id, user, message, logger): - if user: - return MessageHelp(user, BASIC_COMMANDS) - return MessageNotRegistered() - - def handle_message(self, chat_id, user, message, logger): - state, chat_metadata = self.chat_states.get(chat_id) - - # Unregistered user - if not user: - return MessageNotRegistered() - - normalized_message = message.strip().lower() - - if state is None: - # Default state - if normalized_message == COMMAND_WHO.text: - space_status = self.aggregator.get_space_state_for_json(logger) - return MessageWho(user, space_status) - elif normalized_message == COMMAND_HELP.text: - return MessageHelp(user, BASIC_COMMANDS) - elif normalized_message == COMMAND_OUT.text: - return self._handle_checkout(chat_id, user, logger) - elif normalized_message == COMMAND_CHECKIN.text: - self.aggregator.user_entered_space(user.user_id, logger) - space_status = self.aggregator.get_space_state_for_json(logger) - return MessageWho(user, space_status) - else: - return MessageUnknown(user, [COMMAND_WHO.text, COMMAND_OUT.text]) - - elif state == STATE_CONFIRM_CHECKOUT: - if normalized_message == COMMAND_YES.text: - self.aggregator.user_left_space(user.user_id, logger) - self.chat_states.clear(chat_id) - return MessageConfirmedCheckout(user) - elif normalized_message == COMMAND_NO.text: - self.chat_states.clear(chat_id) - return MessageCancelAction() - else: - return MessageUnknown(user, [COMMAND_YES.text, COMMAND_NO.text]) - - elif state == STATE_CONFIRM_VOLUNTEERING: - if normalized_message == COMMAND_YES.text: - registered = self.aggregator.user_volunteers_for_event( - chat_metadata["user_id"], chat_metadata["event"], logger - ) - self.chat_states.clear(chat_id) - return ( - MessageConfirmedVolunteering() - if registered - else MessageVolunteeringNotNecessary() - ) - elif normalized_message == COMMAND_NO.text: - self.chat_states.clear(chat_id) - return MessageCancelAction() - else: - return MessageUnknown(user, [COMMAND_YES.text, COMMAND_NO.text]) - - else: - # Unknown state - should never be here - self.chat_states.clear(chat_id) - return MessageUnknown(user) - - def _handle_checkout(self, chat_id, user, logger): - is_in_space, ts_checkin = self.aggregator.is_user_id_in_space( - user.user_id, logger - ) - if not is_in_space: - return MessageUserNotInSpace(user) - else: - self.chat_states.set(chat_id, STATE_CONFIRM_CHECKOUT) - return MessageConfirmCheckout(user, ts_checkin) - - -class ChatStates(object): - def __init__(self, clock): - self.clock = clock - self.states = {} - - def get(self, chat_id): - value = self.states.get(chat_id) - if not value: - return None, None - state, expiration_ts, metadata = value - if not expiration_ts: - return state, metadata - if self.clock.now() < expiration_ts: - return state, metadata - else: - self.clear(chat_id) - return None, metadata - - def set(self, chat_id, state, expiration_in_min=None, metadata=None): - expiration_ts = None - if expiration_in_min: - expiration_ts = self.clock.now().add(expiration_in_min, "minutes") - self.states[chat_id] = (state, expiration_ts, metadata) - - def clear(self, chat_id): - self.states[chat_id] = None diff --git a/src/aggregator/bots/signal_bot.py b/src/aggregator/bots/signal_bot.py deleted file mode 100644 index c50a68a..0000000 --- a/src/aggregator/bots/signal_bot.py +++ /dev/null @@ -1,76 +0,0 @@ -from functools import partial - -import ravel - -BUS_NAME = "org.asamk.Signal" -PATH_NAME = "/org/asamk/Signal" -IFACE_NAME = "org.asamk.Signal" -SIGNAL_NAME = "MessageReceived" -IN_SIGNATURE = "xsaysas" - - -class SignalBot(object): - def __init__(self, worker_input_queue, aggregator, logger, asyncio_loop): - self.logger = logger.getLogger(subsystem="signal_bot") - self.worker_input_queue = worker_input_queue - self.aggregator = aggregator - self.aggregator.signal_bot = self - self.bus = ravel.system_bus() - self.bus.attach_asyncio(asyncio_loop) - - # Called from the working thread - def send_notification(self, user, notification, logger): - logger = logger.getLogger(subsystem="signal_bot") - logger.info( - f"Sending notification of type {notification.__class__.__name__} to user {user.user_id} {user.full_name}" - ) - self._send_message(notification, user.phone_number) - chat_id = f"signal-{user.phone_number}" - return chat_id - - def start_bot(self): - self.logger.info("Starting Signal BOT") - - self.bus.listen_signal( - path=PATH_NAME, - interface=IFACE_NAME, - name=SIGNAL_NAME, - func=self.handle_message, - fallback=True, - ) - - @ravel.signal(name=SIGNAL_NAME, in_signature=IN_SIGNATURE) - def handle_message(self, msgid, sender, groupIDs, message, attachments): - try: - phone_number = sender - chat_id = f"signal-{phone_number}" - user = self.worker_input_queue.add_task_with_result_blocking( - partial(self.aggregator.get_user_by_phone_number, phone_number), - self.logger, - ) - reply = self.worker_input_queue.add_task_with_result_blocking( - partial(self.aggregator.handle_bot_message, chat_id, user, message), - self.logger, - ) - if not reply: - self.logger.error("Missing reply from BOT logic") - else: - self._send_message(reply, sender) - except Exception: - self.logger.exception("Unexpected exception in message handler") - - def _send_message(self, message, phone_number): - body = message.get_text() - ep = self.bus[BUS_NAME][PATH_NAME].get_interface(IFACE_NAME) - ep.sendMessage(body, [], [phone_number]) - - def stop_bot(self): - self.logger.info("Stopping Signal BOT") - - self.bus.unlisten_signal( - path=PATH_NAME, - interface=IFACE_NAME, - name=SIGNAL_NAME, - func=self.handle_message, - fallback=True, - ) diff --git a/src/aggregator/bots/telegram_bot.py b/src/aggregator/bots/telegram_bot.py deleted file mode 100644 index b3c2bd8..0000000 --- a/src/aggregator/bots/telegram_bot.py +++ /dev/null @@ -1,125 +0,0 @@ -from functools import partial - -from telegram import ReplyKeyboardMarkup, ReplyKeyboardRemove -from telegram.ext import CommandHandler, Filters, MessageHandler, Updater - - -class TelegramBot(object): - def __init__(self, worker_input_queue, aggregator, logger, api_token): - self.logger = logger.getLogger(subsystem="telegram_bot") - self.worker_input_queue = worker_input_queue - self.aggregator = aggregator - self.updater = Updater(api_token) - self.aggregator.telegram_bot = self - - # Called from the working thread - def send_notification(self, user, notification, logger): - logger = logger.getLogger(subsystem="telegram_bot") - logger.info( - f"Sending notification of type {notification.__class__.__name__} to user {user.user_id} {user.full_name}" - ) - self._send_message(notification, user.telegram_user_id) - chat_id = f"telegram-{user.telegram_user_id}" - return chat_id - - def start_bot(self): - self.logger.info("Starting Telegram BOT") - - dp = self.updater.dispatcher - dp.add_handler(CommandHandler("start", self.handle_message)) - dp.add_handler(MessageHandler(Filters.text, self.handle_message)) - dp.add_error_handler(self._error) - - self.updater.start_polling() - - def handle_message(self, bot, update): - try: - telegram_id = str(update.message.chat_id) - chat_id = f"telegram-{telegram_id}" - message = update.message.text - user = self._get_user_by_telegram_id(telegram_id) - self.logger.info( - f'Received message "{message}" from user {user.full_name if user else ""}' - ) - if message.startswith("/start"): - connection_token = get_connection_token_from_message(message) - if connection_token: - user = self._make_new_telegram_association_for_user( - user, connection_token, telegram_id - ) - reply = self._handle_new_bot_conversation(chat_id, user, message) - else: - reply = self._handle_bot_message(chat_id, user, message) - else: - reply = self._handle_bot_message(chat_id, user, message) - if not reply: - self.logger.error("Missing reply from BOT logic") - else: - self._send_message(reply, telegram_id) - except Exception: - self.logger.exception("Unexpected exception in message handler") - - def _send_message(self, message, telegram_user_id): - markdown = message.get_markdown() - reply_markup = ( - ReplyKeyboardMarkup( - [[command.text for command in message.next_commands]], - one_time_keyboard=True, - ) - if message.next_commands - else ReplyKeyboardRemove() - ) - self.updater.bot.send_message( - telegram_user_id, markdown, reply_markup=reply_markup - ) - - def stop_bot(self): - self.logger.info("Stopping Telegram BOT") - if self.updater.running: - self.updater.stop() - - def _error(self, bot, update, error): - self.logger.error(f'Update "{update}" caused error "{error}"', error) - - # -- Proxied aggregator methods (via worker_input_queue) ---- - - def _make_new_telegram_association_for_user( - self, current_user, connection_token, telegram_id - ): - return self.worker_input_queue.add_task_with_result_blocking( - partial( - self.aggregator.make_new_telegram_association_for_user, - current_user, - connection_token, - telegram_id, - ), - self.logger, - ) - - def _get_user_by_telegram_id(self, telegram_id): - return self.worker_input_queue.add_task_with_result_blocking( - partial(self.aggregator.get_user_by_telegram_id, telegram_id), self.logger - ) - - def _handle_new_bot_conversation(self, chat_id, user, message): - return self.worker_input_queue.add_task_with_result_blocking( - partial( - self.aggregator.handle_new_bot_conversation, chat_id, user, message - ), - self.logger, - ) - - def _handle_bot_message(self, chat_id, user, message): - message = message.lstrip("/") - return self.worker_input_queue.add_task_with_result_blocking( - partial(self.aggregator.handle_bot_message, chat_id, user, message), - self.logger, - ) - - -def get_connection_token_from_message(message): - if message.startswith("/start"): - parts = message.split(" ", 1) - if len(parts) > 1: - authorization_token = parts[1] - return authorization_token diff --git a/src/aggregator/database.py b/src/aggregator/database.py index 9a9075a..37193e6 100644 --- a/src/aggregator/database.py +++ b/src/aggregator/database.py @@ -53,31 +53,3 @@ def get_all_tags(self, logger): """ ) return [Tag(row[0], row[1], User(*row[2:])) for row in mycursor] - - def store_telegram_user_id_for_user_id(self, telegram_user_id, user_id, logger): - logger = logger.getLogger(subsystem="mysql") - logger.info( - f"Registering user {user_id} with Telegram User ID {telegram_user_id}" - ) - with self._connection() as db: - mycursor = db.cursor() - mycursor.execute( - """ - UPDATE members_user SET telegram_user_id = %s WHERE id = %s - """, - (telegram_user_id, user_id), - ) - db.commit() - - def delete_telegram_user_id_for_user_id(self, user_id, logger): - logger = logger.getLogger(subsystem="mysql") - logger.info(f"Clearing Telegram User ID for user {user_id}") - with self._connection() as db: - mycursor = db.cursor() - mycursor.execute( - """ - UPDATE members_user SET telegram_user_id = NULL WHERE id = %s - """, - (user_id,), - ) - db.commit() diff --git a/src/aggregator/http_server.py b/src/aggregator/http_server.py index ce1ca9e..7f0f5ad 100644 --- a/src/aggregator/http_server.py +++ b/src/aggregator/http_server.py @@ -122,28 +122,6 @@ async def telegram_token(): ) return Response(token.encode("utf-8"), mimetype="text/plain") - @app.route("/telegram/disconnect", methods=["POST"]) - @with_basic_auth - async def telegram_disconnect(): - request_body = await request.get_data() - request_payload = json.loads(request_body) - await worker_input_queue.add_task_with_result_future( - partial(aggregator.delete_telegram_id_for_user, request_payload["user_id"]), - request.logger, - ) - return Response("Ok", mimetype="text/plain") - - @app.route("/signal/onboard", methods=["POST"]) - @with_basic_auth - async def signal_onboard(): - request_body = await request.get_data() - request_payload = json.loads(request_body) - await worker_input_queue.add_task_with_result_future( - partial(aggregator.onboard_new_signal_user, request_payload["user_id"]), - request.logger, - ) - return Response("Ok", mimetype="text/plain") - @app.route("/notification/test", methods=["POST"]) @with_basic_auth async def notification_test(): diff --git a/src/aggregator/logic.py b/src/aggregator/logic.py index f55a8cd..ac64726 100644 --- a/src/aggregator/logic.py +++ b/src/aggregator/logic.py @@ -1,13 +1,9 @@ # This is where the main business logic lives. -import random from collections import defaultdict -from .bots.bot_logic import BotLogic from .messages import ( - BASIC_COMMANDS, MachineLeftOnNotification, - MessageHelp, ProblemLightLeftOn, ProblemMachineLeftOnBySomeoneElse, ProblemMachineLeftOnByUser, @@ -44,9 +40,6 @@ def __init__( self.checkin_stale_after_hours = checkin_stale_after_hours self.email_adapter = email_adapter self.task_scheduler = task_scheduler - self.bot_logic = BotLogic(self) - self.telegram_bot = None - self.signal_bot = None self.urls = Urls() def _get_user_by_id(self, user_id, logger): @@ -81,27 +74,6 @@ def _get_all_machines(self, logger): # -------------------------------------------------- - def make_new_telegram_association_for_user( - self, current_user, connection_token, telegram_id, logger - ): - if current_user: - # Clear current association (if there is one) - self.delete_telegram_id_for_user(current_user.user_id, logger) - connecting_user = self.register_user_by_telegram_token( - connection_token, telegram_id, logger - ) - return connecting_user - - def get_user_by_telegram_id(self, telegram_id, logger): - user = self.redis_adapter.get_user_by_telegram_id(telegram_id, logger) - if not user: - all_users = self.database_adapter.get_all_users(logger) - self.redis_adapter.set_users_by_ids(all_users, logger) - filtered_users = [u for u in all_users if u.telegram_user_id == telegram_id] - if len(filtered_users) == 1: - user = filtered_users[0] - return user - def get_user_by_phone_number(self, phone_number, logger): user = self.redis_adapter.get_user_by_phone_number(phone_number, logger) if not user: @@ -322,28 +294,6 @@ def _check_out_stale_user(self, user_id, ts_checkin, elapsed_time_in_hours, logg self.crm_adapter.user_checkout(user_id, logger) def send_user_notification(self, user, notification, logger): - if self.telegram_bot and user.uses_telegram_bot(): - chat_id = None - try: - chat_id = self.telegram_bot.send_notification( - user, notification, logger - ) - except Exception: - logger.exception( - "Unexpected exception when trying to notify user via Telegram" - ) - if chat_id: - notification.set_chat_state(chat_id, self.bot_logic) - if self.signal_bot and user.uses_signal_bot(): - chat_id = None - try: - chat_id = self.signal_bot.send_notification(user, notification, logger) - except Exception: - logger.exception( - "Unexpected exception when trying to notify user via Signal" - ) - if chat_id: - notification.set_chat_state(chat_id, self.bot_logic) if user.uses_email(): try: self.email_adapter.send_email_to_user(user, notification, logger) @@ -392,44 +342,6 @@ def lights(self, room, state, logger): logger = logger.getLogger(subsystem="aggregator") self.redis_adapter.set_lights(room, state, logger) - def create_telegram_connect_token(self, user_id, logger): - logger = logger.getLogger(subsystem="aggregator") - token = str(random.randint(10**20, 10**21)) - self.redis_adapter.set_telegram_token(token, user_id, logger) - return token - - def register_user_by_telegram_token(self, token, telegram_user_id, logger): - logger = logger.getLogger(subsystem="aggregator") - user_id = self.redis_adapter.get_user_id_by_telegram_token(token, logger) - if user_id: - user = self._get_user_by_id(user_id, logger) - logger.info( - f"Registering user {user_id} with Telegram User ID {telegram_user_id}" - ) - self.database_adapter.store_telegram_user_id_for_user_id( - telegram_user_id, user.user_id, logger - ) - return user - - def delete_telegram_id_for_user(self, user_id, logger): - self.database_adapter.delete_telegram_user_id_for_user_id(user_id, logger) - all_users = self.database_adapter.get_all_users(logger) - self.redis_adapter.set_users_by_ids(all_users, logger) - - def handle_bot_message(self, chat_id, user, message, logger): - return self.bot_logic.handle_message(chat_id, user, message, logger) - - def handle_new_bot_conversation(self, chat_id, user, message, logger): - return self.bot_logic.handle_new_conversation(chat_id, user, message, logger) - - def onboard_new_signal_user(self, user_id, logger): - logger = logger.getLogger(subsystem="aggregator") - user = self._get_user_by_id(user_id, logger) - if self.signal_bot: - self.signal_bot.send_notification( - user, MessageHelp(user, BASIC_COMMANDS), logger - ) - def send_notification_test(self, user_id, logger): logger = logger.getLogger(subsystem="aggregator") user = self._get_user_by_id(user_id, logger) diff --git a/src/aggregator/logic_tests.py b/src/aggregator/logic_tests.py index a76a6b4..3443bef 100644 --- a/src/aggregator/logic_tests.py +++ b/src/aggregator/logic_tests.py @@ -63,11 +63,7 @@ def test_warn_user_when_he_leaves_and_his_machine_is_on(self): self.aggregator.user_left_space(STEFANO.user_id, self.logger) - self.assertEqual(self.bot_messages, [(1, "ProblemsLeavingSpaceNotification")]) - problems = [ - p.__class__.__name__ for p in self.bot_notification_objects[0].problems - ] - self.assertEqual(problems, ["ProblemMachineLeftOnByUser"]) + self.assertEqual(self.bot_messages, []) def test_warn_user_when_he_leaves_and_another_machine_is_on(self): self.aggregator.user_entered_space(STEFANO.user_id, self.logger) @@ -79,11 +75,7 @@ def test_warn_user_when_he_leaves_and_another_machine_is_on(self): self.aggregator.user_left_space(STEFANO.user_id, self.logger) - self.assertEqual(self.bot_messages, [(1, "ProblemsLeavingSpaceNotification")]) - problems = [ - p.__class__.__name__ for p in self.bot_notification_objects[0].problems - ] - self.assertEqual(problems, ["ProblemMachineLeftOnBySomeoneElse"]) + self.assertEqual(self.bot_messages, []) def test_warn_user_when_he_leaves_lights_on(self): self.aggregator.user_entered_space(STEFANO.user_id, self.logger) @@ -97,11 +89,7 @@ def test_warn_user_when_he_leaves_lights_on(self): self.aggregator.user_left_space(STEFANO.user_id, self.logger) - self.assertEqual(self.bot_messages, [(1, "ProblemsLeavingSpaceNotification")]) - problems = [ - p.__class__.__name__ for p in self.bot_notification_objects[0].problems - ] - self.assertEqual(problems, ["ProblemLightLeftOn"]) + self.assertEqual(self.bot_messages, []) def test_machine_on_and_off(self): self.aggregator.user_entered_space(STEFANO.user_id, self.logger) diff --git a/src/aggregator/main.py b/src/aggregator/main.py index fef648a..72329c0 100644 --- a/src/aggregator/main.py +++ b/src/aggregator/main.py @@ -45,25 +45,6 @@ def _main(config): logger, logging_handler = configure_logging(**config.get("logging", {})) logger.info("Initializing Aggregator service") - # From https://stackoverflow.com/questions/2549939/get-signal-names-from-numbers-in-python - # _signames = {v: k - # for k, v in reversed(sorted(vars(signal).items())) - # if k.startswith('SIG') and not k.startswith('SIG_')} - # - # - # def get_signal_name(signum): - # """Returns the signal name of the given signal number.""" - # return _signames[signum] - # - # # Properly detect Ctrl+C - # def signal_handler(signum, frame): - # print('AAA') - # # print('Received signal {} ({}), stopping...'.format(signum, get_signal_name(signum))) - # # os._exit(1) - # signal.signal(signal.SIGINT, signal_handler) - # signal.signal(signal.SIGTERM, signal_handler) - # signal.signal(signal.SIGABRT, signal_handler) - # Initialize AsyncIO loop = asyncio.get_event_loop() @@ -108,30 +89,6 @@ def _main(config): worker = Worker(worker_input_queue) worker.start_working_in_background_thread() - # Start Telegram BOT - telegram_bot = None - if config.get("telegram_bot"): - try: - from aggregator.bots.telegram_bot import TelegramBot - - telegram_bot = TelegramBot( - worker_input_queue, aggregator, logger, **config["telegram_bot"] - ) - telegram_bot.start_bot() - except Exception: - logger.exception("Unexpected error while starting Telegram BOT") - - # Start Signal BOT - signal_bot = None - if config.get("signal_bot"): - try: - from aggregator.bots.signal_bot import SignalBot - - signal_bot = SignalBot(worker_input_queue, aggregator, logger, loop) - signal_bot.start_bot() - except Exception: - logger.exception("Unexpected error while starting Signal BOT") - # Start cronjobs if "check_stale_checkins" in config: start_checking_for_stale_checkins( @@ -155,8 +112,4 @@ def _main(config): ) # Quit the application - if telegram_bot: - telegram_bot.stop_bot() - if signal_bot: - signal_bot.stop_bot() mqtt_listener_client.stop() diff --git a/src/aggregator/model.py b/src/aggregator/model.py index 8fb4b80..cda5adc 100644 --- a/src/aggregator/model.py +++ b/src/aggregator/model.py @@ -8,7 +8,7 @@ class User( namedtuple( "User", - "user_id first_name last_name email telegram_user_id phone_number uses_signal always_uses_email", + "user_id first_name last_name email phone_number always_uses_email", ) ): @property @@ -18,22 +18,12 @@ def full_name(self): def for_json(self): d = dict(self._asdict()) d["full_name"] = self.full_name - del d["telegram_user_id"] del d["phone_number"] - del d["uses_signal"] del d["always_uses_email"] return d - def uses_telegram_bot(self): - return bool(self.telegram_user_id) - - def uses_signal_bot(self): - return self.uses_signal and self.phone_number - def uses_email(self): - return ( - not self.uses_telegram_bot() and not self.uses_signal_bot() - ) or self.always_uses_email + return self.always_uses_email Tag = namedtuple("Tag", "tag_id tag user") diff --git a/src/aggregator/redis.py b/src/aggregator/redis.py index fd3ee8f..f602fe8 100644 --- a/src/aggregator/redis.py +++ b/src/aggregator/redis.py @@ -20,7 +20,6 @@ def __init__( key_prefix, users_expiration_time_in_sec, pending_machine_activation_timeout_in_sec, - telegram_token_expiration_in_sec, machine_state_timeout_in_minutes, history_lines_expiration_in_days, ): @@ -34,7 +33,6 @@ def __init__( self.pending_machine_activation_timeout_in_sec = ( pending_machine_activation_timeout_in_sec ) - self.telegram_token_expiration_in_sec = telegram_token_expiration_in_sec self.machine_state_timeout_in_minutes = machine_state_timeout_in_minutes self.history_lines_expiration_in_days = history_lines_expiration_in_days @@ -91,19 +89,6 @@ def set_users_by_ids(self, users, logger): ) self.redis.delete(self._k_users_by_telegram_id()) - telegram_users = [u for u in users if u.telegram_user_id] - if len(telegram_users) > 0: - self.redis.hmset( - self._k_users_by_telegram_id(), - dict( - (user.telegram_user_id, json.dumps(user._asdict())) - for user in telegram_users - ), - ) - self.redis.pexpire( - self._k_users_by_telegram_id(), - self.users_expiration_time_in_sec * 1000, - ) self.redis.delete(self._k_users_by_phone_number()) users_with_phone_number = [u for u in users if u.phone_number] @@ -239,28 +224,6 @@ def get_lights_on(self, logger): logger.info("Reading lights ON") return [m.decode("utf-8") for m in self.redis.smembers(self._k_lights_on())] - def set_telegram_token(self, token, user_id, logger): - logger = logger.getLogger(subsystem="redis") - logger.info(f"Set Telegram token for user {user_id}") - self.redis.setex( - self._k_telegram_token(token), - self.telegram_token_expiration_in_sec, - str(user_id), - ) - - def get_user_id_by_telegram_token(self, token, logger): - logger = logger.getLogger(subsystem="redis") - logger.info(f"Get Telegram token {token}") - value = self.redis.get(self._k_telegram_token(token)) - return int(value) if value else None - - def get_user_by_telegram_id(self, telegram_id, logger): - logger = logger.getLogger(subsystem="redis") - logger.info(f"Get user with Telegram ID {telegram_id}") - data = self.redis.hget(self._k_users_by_telegram_id(), telegram_id) - if data: - return User(**json.loads(data)) - def get_user_by_phone_number(self, phone_number, logger): logger = logger.getLogger(subsystem="redis") logger.info(f"Get user with phone number {phone_number}") @@ -330,9 +293,6 @@ def _k_nudge(self, nudge_key): def _k_space_open(self): return f"{self.key_prefix}:so" - def _k_telegram_token(self, token): - return f"{self.key_prefix}:tt{token}" - def _k_users_by_id(self): return f"{self.key_prefix}:ui" diff --git a/src/aggregator/testing_utils.py b/src/aggregator/testing_utils.py index a68cfc9..be01da3 100644 --- a/src/aggregator/testing_utils.py +++ b/src/aggregator/testing_utils.py @@ -14,9 +14,7 @@ first_name="Stefano", last_name="Masini", email="stefano@stefanomasini.com", - telegram_user_id="1234", phone_number="+316123456", - uses_signal=True, always_uses_email=True, ) @@ -25,9 +23,7 @@ first_name="Bob", last_name="de Bouwer", email="bob@bouwer.com", - telegram_user_id="2345", phone_number="+316456789", - uses_signal=True, always_uses_email=True, ) @@ -86,7 +82,6 @@ def setUp(self): 60, 90, 60, - 60, 7, ) self.crm_adapter = MockCrmAdapter() @@ -103,7 +98,6 @@ def setUp(self): self.task_scheduler, 5, ) - self.aggregator.signal_bot = self self.bot_messages = [] self.emails_sent = [] self.bot_notification_objects = [] @@ -132,12 +126,3 @@ def send_email_to_user(self, user, message, logger): def send_email(self, name, email, message, logger): self.emails_sent.append((name, email, message.__class__.__name__)) - - def send_bot_message(self, user, message): - reply = self.aggregator.handle_bot_message( - f"signal-{user.phone_number}", user, message, self.logger - ) - if not reply: - raise Exception("Missing reply from BOT logic") - else: - self.send_notification(user, reply, self.logger)