Skip to content

Commit b99b910

Browse files
authored
feat: Decouple Telethon interactions from API to reduce the number of active sessions (#185)
2 parents 2ff4cf5 + c3dbb97 commit b99b910

27 files changed

Lines changed: 663 additions & 519 deletions

‎backend/api/routes/admin/chat/rule/whitelist.py‎

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -62,13 +62,12 @@ async def update_chat_whitelist_rule(
6262
requestor=request.state.user,
6363
chat_slug=slug,
6464
)
65-
action.update(
65+
result = action.update(
6666
rule_id=rule_id,
6767
name=rule.name,
6868
description=rule.description,
6969
is_enabled=rule.is_enabled,
7070
)
71-
result = await action.set_content(rule_id=rule_id, content=rule.users)
7271
return WhitelistRuleFDO.model_validate(result.model_dump())
7372

7473

‎backend/api/routes/user.py‎

Lines changed: 0 additions & 41 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,5 @@
1-
from typing import Annotated
2-
31
from fastapi import APIRouter, Depends
42
from fastapi.exceptions import HTTPException
5-
from fastapi.params import Query
63
from sqlalchemy.orm import Session
74
from starlette.requests import Request
85
from starlette.status import HTTP_400_BAD_REQUEST, HTTP_200_OK
@@ -137,41 +134,3 @@ async def set_user_wallet(
137134
# No need to refresh wallet details if it is already a tracked wallet
138135
task_id=None,
139136
)
140-
141-
142-
@user_router.delete(
143-
"/wallet",
144-
name="Disconnect user wallet from the chat",
145-
description="Disconnect wallet from the chat.",
146-
tags=["Wallet"],
147-
responses={
148-
HTTP_200_OK: {"model": UserFDO},
149-
HTTP_400_BAD_REQUEST: {
150-
"description": "Chat not found",
151-
"model": BaseExceptionFDO,
152-
},
153-
},
154-
)
155-
async def disconnect_wallet(
156-
request: Request,
157-
chat_slug: Annotated[
158-
str,
159-
Query(
160-
...,
161-
alias="chatSlug",
162-
description="Chat slug for which wallet will be disconnected",
163-
),
164-
],
165-
db_session: Session = Depends(get_db_session),
166-
) -> UserFDO:
167-
wallet_action = WalletAction(db_session)
168-
try:
169-
await wallet_action.disconnect_wallet(
170-
user_id=request.state.user.id, chat_slug=chat_slug
171-
)
172-
except TelegramChatNotExists:
173-
raise HTTPException(
174-
detail="Chat not found",
175-
status_code=HTTP_400_BAD_REQUEST,
176-
)
177-
return UserFDO.from_orm(request.state.user)

‎backend/community_manager/actions/chat.py‎

Lines changed: 230 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,8 @@
1010
HideRequesterMissingError,
1111
RPCError,
1212
UserIsBlockedError,
13+
UserAdminInvalidError,
14+
ChatAdminRequiredError,
1315
)
1416
from telethon.utils import get_peer_id
1517

@@ -33,11 +35,13 @@
3335
TelegramChatPublicError,
3436
TelegramChatAlreadyExists,
3537
)
36-
from core.models.chat import TelegramChat
38+
from core.exceptions.telethon import MissingChatEntityError, MissingUserEntityError
39+
from core.models.chat import TelegramChat, TelegramChatUser
3740
from core.services.cdn import CDNService
3841
from core.services.chat import TelegramChatService
3942
from core.services.chat.rule.gift import TelegramChatGiftCollectionService
4043
from core.services.chat.rule.sticker import TelegramChatStickerCollectionService
44+
from core.services.chat.rule.whitelist import TelegramChatExternalSourceService
4145
from core.services.chat.user import TelegramChatUserService
4246
from core.services.superredis import RedisService
4347
from core.services.supertelethon import ChatPeerType, TelethonService
@@ -55,12 +59,11 @@ def __init__(
5559
self.telegram_chat_user_service = TelegramChatUserService(db_session)
5660
self.redis_service = RedisService()
5761
self.cdn_service = CDNService()
58-
self.authorization_action = AuthorizationAction(
59-
db_session, telethon_client=telethon_client
60-
)
62+
self.authorization_action = AuthorizationAction(db_session)
6163
self.telethon_service = TelethonService(
6264
client=telethon_client,
6365
bot_token=community_manager_settings.telegram_bot_token,
66+
session_path=community_manager_settings.telegram_session_path,
6467
)
6568

6669
async def _get_chat_data(
@@ -680,6 +683,64 @@ async def on_join_request(
680683
},
681684
)
682685

686+
async def enable(self, chat_id: int) -> TelegramChat:
687+
"""
688+
This method will enable the chat by setting the invite link and updating status in the DB
689+
690+
:param chat_id: The unique identifier of the Telegram chat to enable.
691+
"""
692+
chat = self.telegram_chat_service.get(chat_id)
693+
if chat.is_enabled:
694+
logger.debug(
695+
f"Chat {chat.id!r} is already enabled. Skipping enable operation..."
696+
)
697+
return chat
698+
699+
await self.telethon_service.start()
700+
try:
701+
peer = await self.telethon_service.get_chat(entity=chat.id)
702+
invite_link = await self.telethon_service.get_invite_link(chat=peer)
703+
chat = self.telegram_chat_service.refresh_invite_link(
704+
chat_id=chat.id, invite_link=invite_link.link
705+
)
706+
logger.info(
707+
f"Updated invite link of chat {chat.id!r} to {invite_link.link!r} and enabled it."
708+
)
709+
except ChatAdminRequiredError:
710+
logger.exception(f"Insufficient privileges to enable chat {chat.id!r}")
711+
raise
712+
except RPCError:
713+
logger.exception(f"Failed to enable chat {chat.id!r}")
714+
raise
715+
finally:
716+
await self.telethon_service.stop()
717+
718+
return chat
719+
720+
async def disable(self, chat_id: int) -> TelegramChat:
721+
"""
722+
This method will disable the chat by setting the invite link and updating status in the DB
723+
:param chat_id: The unique identifier of the Telegram chat to disable.
724+
"""
725+
chat = self.telegram_chat_service.get(chat_id)
726+
await self.telethon_service.start()
727+
try:
728+
await self.telethon_service.revoke_chat_invite(
729+
chat_id=chat.id, link=chat.invite_link
730+
)
731+
chat = self.telegram_chat_service.disable(chat)
732+
logger.info(f"Removed invite link of chat {chat.id!r} and disabled it.")
733+
except ChatAdminRequiredError:
734+
logger.error(f"Insufficient privileges to disable chat {chat.id!r}")
735+
raise
736+
except RPCError:
737+
logger.exception(f"Failed to disable chat {chat.id!r}")
738+
raise
739+
finally:
740+
await self.telethon_service.stop()
741+
742+
return chat
743+
683744

684745
class CommunityManagerTaskChatAction:
685746
def __init__(self, db_session: Session):
@@ -812,10 +873,11 @@ async def sanity_chat_checks(self, telethon_client: TelegramClient) -> None:
812873
else:
813874
logger.info(f"Found {len(chat_members)} chat members to validate")
814875

815-
authorization_action = AuthorizationAction(
816-
self.db_session, telethon_client=telethon_client
876+
community_user_action = CommunityManagerUserChatAction(
877+
db_session=self.db_session,
878+
telethon_client=telethon_client,
817879
)
818-
await authorization_action.kick_ineligible_chat_members(
880+
await community_user_action.kick_ineligible_chat_members(
819881
chat_members=chat_members
820882
)
821883
logger.info(
@@ -844,3 +906,164 @@ def fallback_update_chat_members(self, dto: TargetChatMembersDTO) -> None:
844906
self.redis_service.add_to_set(
845907
UPDATED_STICKERS_USER_IDS, *map(str, dto.sticker_owners_ids)
846908
)
909+
910+
async def refresh_enabled(self, telethon_client: TelegramClient) -> None:
911+
"""
912+
Refreshes all enabled Telegram chat external sources.
913+
914+
This method retrieves the list of enabled external sources and performs validation
915+
to refresh them.
916+
For removed members, it handles appropriate actions such as kicking ineligible chat members.
917+
This ensures synchronization between the source's metadata and the chat's current state.
918+
"""
919+
telegram_chat_external_source_service = TelegramChatExternalSourceService(
920+
self.db_session
921+
)
922+
sources = telegram_chat_external_source_service.get_all(enabled_only=True)
923+
community_user_action = CommunityManagerUserChatAction(
924+
db_session=self.db_session, telethon_client=telethon_client
925+
)
926+
for source in sources:
927+
logger.info(
928+
f"Refreshing enabled chat source {source.chat_id!r} for chat {source.chat_id!r} with URL {source.url!r}"
929+
)
930+
# It should not raise, but log any validation error and continue
931+
diff = await telegram_chat_external_source_service.validate_external_source(
932+
source, raise_for_error=False
933+
)
934+
if not diff:
935+
logger.warning(f"Validation for {source.url!r} failed. Continue...")
936+
continue
937+
938+
if diff.removed:
939+
chat_members = self.telegram_chat_user_service.get_all(
940+
user_ids=diff.removed
941+
)
942+
await community_user_action.kick_ineligible_chat_members(
943+
chat_members=chat_members
944+
)
945+
# Set content only after the source was refreshed to ensure
946+
# no new attempts to kick users that are already kicked will be made
947+
telegram_chat_external_source_service.set_content(source, diff.current)
948+
949+
logger.info("All enabled chat sources refreshed.")
950+
951+
952+
class CommunityManagerUserChatAction:
953+
def __init__(
954+
self, db_session: Session, telethon_client: TelegramClient | None = None
955+
):
956+
self.db_session = db_session
957+
self.telegram_chat_user_service = TelegramChatUserService(db_session)
958+
self.authorization_action = AuthorizationAction(db_session)
959+
self.telethon_service = TelethonService(
960+
client=telethon_client,
961+
bot_token=community_manager_settings.telegram_bot_token,
962+
)
963+
964+
async def kick_chat_member(self, chat_member: TelegramChatUser) -> None:
965+
"""
966+
Kicks a specified chat member from the chat. It ensures that the bot
967+
has sufficient privileges to perform the action and sends a notification
968+
to the user if they allow direct messages. The method handles exceptions
969+
arising due to administrative restrictions or RPC errors and logs
970+
appropriate messages for each case.
971+
972+
:param chat_member: A TelegramChatUser object representing the user to be
973+
kicked from the chat. Must be a bot-managed user with attributes defining
974+
their chat, user ID, and permission states.
975+
"""
976+
if not chat_member.is_managed:
977+
logger.warning(
978+
f"Attempt to kick non-managed chat member {chat_member.chat_id=} and {chat_member.user_id=}. Skipping."
979+
)
980+
return
981+
982+
if chat_member.chat.insufficient_privileges:
983+
logger.warning(
984+
f"Attempt to kick chat member {chat_member.chat_id=} and {chat_member.user_id=} "
985+
f"failed as bot was lacking privileges to manage the chat. Skipping."
986+
)
987+
return
988+
989+
await self.telethon_service.start()
990+
try:
991+
await self.telethon_service.kick_chat_member(
992+
chat_id=chat_member.chat_id,
993+
telegram_user_id=chat_member.user.telegram_id,
994+
)
995+
if chat_member.user.allows_write_to_pm:
996+
try:
997+
await self.telethon_service.send_message(
998+
chat_id=chat_member.user.telegram_id,
999+
message=f"You were kicked out of the **{chat_member.chat.title}**.",
1000+
)
1001+
except RPCError as e:
1002+
logger.error(
1003+
f"Failed to send message to user {chat_member.user.telegram_id!r} "
1004+
f"while kicking them from chat {chat_member.chat_id!r}",
1005+
exc_info=e,
1006+
)
1007+
self.telegram_chat_user_service.delete(
1008+
chat_id=chat_member.chat_id, user_id=chat_member.user.id
1009+
)
1010+
logger.info(
1011+
f"User {chat_member.user.telegram_id!r} was kicked from chat {chat_member.chat_id!r}"
1012+
)
1013+
except UserAdminInvalidError as e:
1014+
logger.warning(
1015+
f"Failed to kick user {chat_member.user.telegram_id!r} from chat {chat_member.chat_id!r} as bot user lacks admin privileges",
1016+
exc_info=e,
1017+
)
1018+
telegram_chat_service = TelegramChatService(self.db_session)
1019+
telegram_chat_service.set_insufficient_privileges(
1020+
chat_id=chat_member.chat_id, value=True
1021+
)
1022+
logger.info(
1023+
f"Set insufficient privileges flag for chat {chat_member.chat_id!r}."
1024+
)
1025+
except RPCError as e:
1026+
logger.error(
1027+
f"Failed to kick user {chat_member.user.telegram_id!r} from chat {chat_member.chat_id!r}",
1028+
exc_info=e,
1029+
)
1030+
finally:
1031+
await self.telethon_service.stop()
1032+
1033+
async def kick_ineligible_chat_members(
1034+
self,
1035+
chat_members: list[TelegramChatUser],
1036+
) -> None:
1037+
"""
1038+
Kicks ineligible chat members from a chat group asynchronously. The method checks
1039+
the eligibility of chat members provided and attempts to remove members deemed
1040+
ineligible. Logging is performed to document successful removals and capture any
1041+
exceptions encountered while processing.
1042+
1043+
:param chat_members: List of chat members to be evaluated and potentially removed.
1044+
:return: This function does not return any value.
1045+
:raises MissingChatEntityError: Raised when the chat entity is missing for a member.
1046+
:raises MissingUserEntityError: Raised when the user entity is missing for a member.
1047+
"""
1048+
ineligible_members = self.authorization_action.get_ineligible_chat_members(
1049+
chat_members=chat_members
1050+
)
1051+
if not ineligible_members:
1052+
logger.info("No ineligible chat members found")
1053+
return
1054+
1055+
logger.info(f"Found {len(ineligible_members)} ineligible chat members")
1056+
1057+
for member in ineligible_members:
1058+
try:
1059+
await self.kick_chat_member(member)
1060+
except MissingChatEntityError as e:
1061+
logger.error(
1062+
f"Failed to kick chat member {member.chat_id=} and {member.user_id=} as chat entity is missing",
1063+
exc_info=e,
1064+
)
1065+
except MissingUserEntityError as e:
1066+
logger.error(
1067+
f"Failed to kick chat member {member.chat_id=} and {member.user_id=} as user entity is missing",
1068+
exc_info=e,
1069+
)

‎backend/community_manager/entrypoint.py‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,9 @@
1818
ChatAdminChangeEventBuilder,
1919
)
2020

21-
logging.basicConfig(level=logging.DEBUG)
21+
logging.basicConfig(
22+
format="%(asctime)s - %(name)s - %(levelname)s - %(message)s", level=logging.DEBUG
23+
)
2224
logger = logging.getLogger(__name__)
2325

2426

0 commit comments

Comments
 (0)