Stabilize staff concurrency and premium emoji; enable verified Drone deployment
Some checks failed
continuous-integration/drone/push Build is failing
Some checks failed
continuous-integration/drone/push Build is failing
This commit is contained in:
97
tests/test_chat_delivery.py
Normal file
97
tests/test_chat_delivery.py
Normal file
@@ -0,0 +1,97 @@
|
||||
import asyncio
|
||||
from datetime import datetime, timezone
|
||||
from types import SimpleNamespace
|
||||
from unittest.mock import AsyncMock
|
||||
|
||||
from aiogram import Bot
|
||||
from aiogram.exceptions import TelegramRetryAfter
|
||||
from aiogram.methods import CopyMessage
|
||||
from aiogram.types import Message, Chat, User as TelegramUser
|
||||
from sqlalchemy import select
|
||||
|
||||
from src.core.broadcast_services import BroadcastService
|
||||
from src.core.chat_services import ChatMessageService
|
||||
from src.core.database import async_session_maker
|
||||
from src.core.models import ChatMessage, BlockedUser
|
||||
from src.core.services import UserService
|
||||
from src.handlers.chat_handlers import _is_message_processed, _processed_messages, _copy_with_sender
|
||||
from test_dispatcher import TelegramStub
|
||||
|
||||
|
||||
def test_equal_message_numbers_from_different_users_are_independent():
|
||||
_processed_messages.clear()
|
||||
assert not _is_message_processed(101, 7)
|
||||
assert not _is_message_processed(102, 7)
|
||||
assert _is_message_processed(101, 7)
|
||||
|
||||
|
||||
async def test_moderation_lookup_includes_receiving_chat():
|
||||
async with async_session_maker() as session:
|
||||
first = await UserService.get_or_create_user(session, 101)
|
||||
second = await UserService.get_or_create_user(session, 102)
|
||||
one = ChatMessage(user_id=first.id, telegram_message_id=7, message_type="text",
|
||||
forwarded_message_ids={"900001": 41})
|
||||
two = ChatMessage(user_id=second.id, telegram_message_id=7, message_type="text",
|
||||
forwarded_message_ids={"900002": 41})
|
||||
session.add_all([one, two])
|
||||
await session.commit()
|
||||
assert (await ChatMessageService.get_message_by_telegram_id(session, 7, chat_id=101)).id == one.id
|
||||
assert (await ChatMessageService.get_message_by_telegram_id(session, 41, chat_id=900002)).id == two.id
|
||||
assert await ChatMessageService.get_message_by_telegram_id(session, 41, chat_id=123) is None
|
||||
assert await ChatMessageService.get_message_by_telegram_id(session, 7) is None
|
||||
|
||||
|
||||
async def test_chat_preserves_literal_sender_and_text_and_keeps_max_length_messages():
|
||||
stub = TelegramStub()
|
||||
bot = Bot("123456:TEST_TOKEN_FOR_ISOLATED_TESTS", session=stub)
|
||||
msg = Message(message_id=2, date=datetime.now(timezone.utc), chat=Chat(id=101, type="private"),
|
||||
from_user=TelegramUser(id=101, first_name="<name>", is_bot=False), text="<b>literal</b> & value").as_(bot)
|
||||
assert await _copy_with_sender(msg, 102, "<name>") == 100
|
||||
assert "<name>" in stub.calls[0].text
|
||||
assert "<b>literal</b> & value" in stub.calls[0].text
|
||||
assert stub.calls[0].parse_mode is None
|
||||
assert len(stub.calls[0].entities) == 1 # Only the explicit sender header is formatted.
|
||||
stub.calls.clear()
|
||||
long_message = msg.model_copy(update={"text": "x" * 4096}).as_(bot)
|
||||
await _copy_with_sender(long_message, 102, "Sender")
|
||||
assert [call.__api_method__ for call in stub.calls] == ["sendMessage", "copyMessage"]
|
||||
|
||||
|
||||
async def test_parallel_broadcast_retries_do_not_deadlock_delivery_slots():
|
||||
service = BroadcastService()
|
||||
async def send(index):
|
||||
message = SimpleNamespace(copy_to=AsyncMock(side_effect=[
|
||||
TelegramRetryAfter(method=CopyMessage(chat_id=index, from_chat_id=1, message_id=1),
|
||||
message="Retry", retry_after=0), SimpleNamespace(message_id=1)]))
|
||||
result = await service.send_message_to_user(None, SimpleNamespace(telegram_id=index), message)
|
||||
assert message.copy_to.await_count == 2
|
||||
return result
|
||||
results = await asyncio.wait_for(asyncio.gather(*(send(index) for index in range(8))), timeout=8)
|
||||
assert all(success for success, _ in results)
|
||||
|
||||
|
||||
async def test_parallel_delivery_failures_create_one_blocked_user_record():
|
||||
service = BroadcastService()
|
||||
async def mark():
|
||||
async with async_session_maker() as session:
|
||||
await service.mark_user_blocked(session, 200, "blocked_bot", "test")
|
||||
await asyncio.gather(*(mark() for _ in range(8)))
|
||||
async with async_session_maker() as session:
|
||||
record = (await session.scalars(select(BlockedUser))).one()
|
||||
assert record.telegram_id == 200
|
||||
assert record.attempt_count == 8
|
||||
|
||||
|
||||
async def test_inactivity_job_preserves_existing_delivery_block_and_continues():
|
||||
from datetime import timedelta
|
||||
from src.core.activity_service import ActivityService
|
||||
async with async_session_maker() as session:
|
||||
for number in (201, 202):
|
||||
user = await UserService.get_or_create_user(session, number)
|
||||
user.is_registered = True
|
||||
user.last_activity = datetime.now(timezone.utc) - timedelta(days=31)
|
||||
await session.commit()
|
||||
await BroadcastService().mark_user_blocked(session, 201, "blocked_bot", "test")
|
||||
assert await ActivityService.mark_inactive_users(session) == 1
|
||||
assert await ActivityService.mark_inactive_users(session) == 0
|
||||
assert await session.scalar(select(BlockedUser.error_type).where(BlockedUser.telegram_id == 201)) == "blocked_bot"
|
||||
Reference in New Issue
Block a user