Как я написал свой SMTP‑сервер, чтобы не пропускать сообщения заказчиков с фриланс‑бирж

от автора

Я фрилансер, и я ненавижу мониторить биржи. FL, Kwork, Freelance, куча Telegram‑каналов — каждое утро одна и та же рутина. Открываешь вкладку, обновляешь, листаешь, понимаешь что всё мимо, переходишь на следующую. Кольцо.

Подсчитал как‑то — два‑три часа в день уходит не на работу. На обновление страниц.

Решил автоматизировать. Написал агрегатор, который собирает заказы с четырёх площадок в единую ленту с фильтрами, присылает уведомления в Telegram и пересылает ответы заказчиков из бирж прямо в чат. Под капотом Python 3.13, FastAPI, Vue 3, PostgreSQL. Всё крутится в одном asyncio event loop.

В этой статье расскажу про самые интересные технические штуки, на которые убил больше всего времени: самописный SMTP‑сервер на чистом asyncio, парсинг JavaScript state data вместо HTML, обход TLS fingerprinting на Kwork и in‑process event bus на asyncio.Condition вместо Redis.

Зачем свой SMTP

Есть четыре биржи. У каждой свои уведомления: когда заказчик пишет в личку, биржа шлёт email. Вопрос — как получить это сообщение в реальном времени и переслать пользователю в Telegram?

Первый вариант очевиден — IMAP polling. Подключаться к почтовому ящику каждые N секунд и проверять новые письма. Проблема в том, что это медленно, неэффективно, а ещё Gmail и Мейл имеют жёсткие лимиты на частоту IMAP‑соединений. Можно конечно IMAP IDLE, но это отдельное TCP‑соединение на каждый ящик, и не все провайдеры его нормально поддерживают.

Второй вариант — принять письмо напрямую. Завести MX‑запись на свой домен и поднять SMTP‑сервер. Пользователь настраивает пересылку с Мейла или Яндекса на адрес u{telegram_id}@inbox.mydomain.ru, и письмо приходит мгновенно.

Я пошёл вторым путём. Но городить полноценный Postfix с очередями, milter‑фильтрами и конфигами на 500 строк ради приёма уведомлений от бирж — перебор. Мне нужен минимальный SMTP‑ресивер: принять письмо, распарсить, понять от какой биржи, вытащить текст сообщения заказчика и дёрнуть Telegram API. Всё.

121 строка и ноль внешних зависимостей

Вместо aiosmtpd (она тянет за собой атрибуты и хуки, которые мне не нужны) я написал сервер с нуля на asyncio.start_server. Весь протокол — это конечный автомат с двумя состояниями: командный режим и режим данных.

class AsyncSMTPServer:    async def handle_client(self, reader, writer):        writer.write(b"220 mail.example.ru ESMTP Inbound Server Ready\r\n")        await writer.drain()        recipient, sender = "", ""        data_mode = False        raw_bytes = bytearray()        while True:            line = await reader.readline()            if not line:                break            if data_mode:                if line in (b".\r\n", b".\n"):                    data_mode = False                    msg = email.message_from_bytes(                        bytes(raw_bytes), policy=default                    )                    # ... парсинг MIME, создание InboundEmail,                    # запуск обработки в фоне                    asyncio.create_task(                        self.process_usecase.execute(inbound_email)                    )                    writer.write(b"250 2.0.0 OK Message accepted\r\n")                    await writer.drain()                    raw_bytes.clear()                else:                    raw_bytes.extend(line)                continue            cmd_upper = line.decode("utf-8", errors="ignore").strip().upper()            if cmd_upper.startswith("EHLO") or cmd_upper.startswith("HELO"):                writer.write(b"250 mail.example.ru OK\r\n")            elif cmd_upper.startswith("MAIL FROM:"):                sender = ...  # парсинг из угловых скобок                writer.write(b"250 2.1.0 Sender OK\r\n")            elif cmd_upper.startswith("RCPT TO:"):                recipient = ...  # адрес вида u1524607402@inbox...                writer.write(b"250 2.1.5 Recipient OK\r\n")            elif cmd_upper == "DATA":                data_mode = True                writer.write(b"354 Start mail input\r\n")            elif cmd_upper == "QUIT":                writer.write(b"221 Bye\r\n")                break            elif cmd_upper in ("RSET", "NOOP"):                writer.write(b"250 OK\r\n")            else:                writer.write(b"500 Command unrecognized\r\n")            await writer.drain()

Вот и весь протокол. HELO, MAIL FROM, RCPT TO, DATA, QUIT, RSET, NOOP — семь команд. В режиме DATA строки складываются в bytearray до строки‑маркера . (точка на отдельной строке, конец тела письма по RFC 5321). Потом стандартный email.message_from_bytes разбирает всё: multipart, charset, quoted‑printable, base64.

Ключевой момент — asyncio.create_task. Обработка письма (поиск юзера в БД, отправка в Telegram) запускается фоновой задачей, не блокируя текущее TCP‑соединение. Если Мейл шлёт два письма подряд, второе не ждёт пока первое допроцессится.

Адресация через email

Telegram ID пользователя зашит прямо в email‑адрес: u1524607402@inbox.mydomain.ru. Домен‑сущность InboundEmail извлекает его регуляркой r"(?:u)?(\d{5,12})@" — 20 строк на весь доменный слой.

@dataclassclass InboundEmail:    recipient: str    sender: str    subject: str    body_text: str    body_html: str    def extract_telegram_id(self) -> Optional[int]:        match = re.search(r"(?:u)?(\d{5,12})@", self.recipient)        return int(match.group(1)) if match else None

Пользователь настраивает пересылку в Мейле или Яндексе один раз, дальше всё работает автоматически. Мейл шлёт подтверждение пересылки — это тоже приходит на SMTP‑сервер, парсится и пересылается юзеру ссылкой в Telegram для быстрого клика.

Парсинг писем: каждая биржа мучает HTML по‑своему

Когда письмо пришло и MIME распарсен, нужно понять: это Kwork, FL, Freelance или подтверждение пересылки от Мейла? И вытащить оттуда текст сообщения заказчика.

Звучит просто. На практике каждая биржа формирует HTML письма по‑своему, и ни одна не думала о том, чтобы это было удобно парсить.

def parse_inbound_email(email: InboundEmail) -> ParsedEmailResult:    # 1. Фильтруем маркетинг: news@, promo@, newsletter@ - мимо    if any(addr in email.sender.lower() for addr in ignored_senders):        return IgnoredEmailResult(reason=f"marketing_sender:{email.sender}")    # 2. Подтверждения пересылки (Mail.ru, Яндекс, Gmail)    if "подтверд" in full_content.lower() or "пересылк" in full_content.lower():        links = re.findall(r'https?://[^\s<>"\']+', full_content)        ...        return MailConfirmationResult(confirmation_link=conf_link, ...)    # 3. Kwork: текст в <i>, ник из "от Username" или kwork.ru/inbox/    if "kwork" in email.sender.lower():        italic_match = re.search(r'<i[^>]*>(.*?)</i>', email.body_html, ...)        user_match = re.search(r'от\s+([a-zA-Z0-9_-]+)', full_content, ...)        ...    # 4. FL.ru: текст между ------ разделителями    if "fl.ru" in email.sender.lower():        dash_match = re.search(r'------\s*(.*?)\s*------', cleaned_html, ...)        ...    # 5. Freelance.ru: текст в <h4>    if "freelance.ru" in email.sender.lower():        h4_match = re.search(r'<h4[^>]*>(.*?)</h4>', email.body_html, ...)        ...

Kwork заворачивает текст сообщения в тег <i>. FL обрамляет содержимое шестью дефисами с каждой стороны ------. Freelance кладёт в <h4>. Ни у одного нет машиночитаемых заголовков или structured data. Чистая эвристика на регулярках.

Если письмо не подпадает ни под одну биржу — просто пересылаем первые 300 символов и тему. Так не теряется ничего.

А зачем вообще тут чистая доменная модель и Union‑тип с четырьмя вариантами? А вот зачем: IgnoredEmailResult позволяет дальше по цепочке не гонять маркетинговый спам, MailConfirmationResult рендерит кнопку с прямой ссылкой, ChatNotificationResult сохраняет сообщение в историю чатов. У каждого типа своя обработка в юзкейсе, и если биржа изменит формат — меняешь один парсер, а не весь пайплайн.

Как я скрепю четыре площадки и почему у Kwork самый хитрый скрапер

FL: RSS и selectolax

FL — самый простой. У них есть открытый RSS‑фид fl.ru/rss/all.xml. Парсим стандартным xml.etree.ElementTree, из каждого <item> достаём ссылку, заголовок, бюджет (он прямо в title в скобках, ага).

Но в RSS описание обрезано. Для полного текста иду на страницу проекта и вытаскиваю div.fl-project-content__description-text через selectolax. Почему selectolax, а не BeautifulSoup? Тупо быстрее. BS4 на тысяче страниц уже начинает тормозить. Selectolax — это биндинг к C‑парсеру Modest/Lexbor, работает в разы шустрее.

async def scrape(self, existing_ids=None):    async with httpx.AsyncClient(follow_redirects=True, timeout=12.0) as client:        await self._scrape_feed(client, "https://www.fl.ru/rss/all.xml", "gigs", ...)        await self._scrape_feed(client, "https://www.fl.ru/rss/office.xml", "vacancies", ...)

Два фида — заказы и вакансии. Всё асинхронно через httpx.

Kwork: curl_cffi и window.stateData

А вот Kwork — это отдельная история. У них Cloudflare WAF, который проверяет TLS fingerprint. Обычный httpx или aiohttp с дефолтными настройками TLS отдают характерный отпечаток, который Cloudflare мгновенно палит.

Решение — curl_cffi. Это Python‑биндинг к libcurl с поддержкой TLS fingerprint impersonation. Один аргумент impersonate="chrome" — и библиотека выставляет cipher suites, extensions, ALPN в точности как настоящий Chrome:

async with AsyncSession() as session:    response = await session.get(        url, headers=self.headers, impersonate="chrome", timeout=15.0    )

Но это полдела. Kwork рендерит список проектов на клиенте через Vue.js. Карточки в HTML — это пустые контейнеры, а данные лежат в JavaScript‑объекте window.stateData, который инжектится в <script> тег.

Вместо того чтобы городить Selenium или Playwright, я вытаскиваю JSON прямо из JavaScript:

def extract_state_data(script_text: str) -> Optional[dict]:    match = re.search(r"(?:window\.)?stateData\s*=\s*(\{)", script_text)    if not match:        return None    start_idx = match.start(1)    brace_count = 0    in_string = False    escape = False    quote_char = None    js_value = []    for i in range(start_idx, len(script_text)):        char = script_text[i]        if escape:            escape = False            js_value.append(char)            continue        if char == "\\":            escape = True            js_value.append(char)            continue        if in_string:            if char == quote_char:                in_string = False            js_value.append(char)            continue        if char in ('"', "'"):            in_string = True            quote_char = char        if char == "{":            brace_count += 1        elif char == "}":            brace_count -= 1        js_value.append(char)        if brace_count == 0:            break    return json.loads("".join(js_value))

Ручной парсер с подсчётом фигурных скобок, обработкой кавычек и escape‑символов. json.loads не сработает без правильного извлечения — объект вложен в произвольный JavaScript, а регулярка stateData\s*=\s*({.*}) жадно захватит лишнее.

Из stateData достаю сразу всё: список проектов (wantsListData.wants), полное дерево категорий, статусы, бюджеты, логины заказчиков, процент найма, прикреплённые файлы ТЗ. Один HTTP‑запрос на страницу, ноль JavaScript‑рендеринга.

Freelance: параллельный парсинг с семафором

Freelance отдаёт HTML нормально, без JS‑рендеринга. Но на странице списка — только краткое описание. Полное описание, данные заказчика, файлы ТЗ — это всё на детальной странице каждого заказа.

Парсить детальные страницы последовательно при 25 заказах на странице — слишком медленно. Параллельно без ограничений — Freelance забанит по IP. Решение — asyncio.Semaphore(5):

semaphore = asyncio.Semaphore(5)  # max 5 параллельных запросов# собираем задачи...detail_tasks = [    self._fetch_task_detail(client, t["detail_url"], t["brief_desc"], semaphore)    for t in page_tasks]details_results = await asyncio.gather(*detail_tasks)

Пять параллельных запросов на детальные страницы, остальные ждут в очереди. С паузой 0.1с между запросами внутри семафора — сервер не жалуется.

Telegram: Telethon userbot в реальном времени

Скрепить Telegram‑каналы обычным HTTP нельзя. Нужен клиент, подписанный на каналы. Я поднял userbot на Telethon, который слушает events.NewMessage и в реальном времени обрабатывает каждый пост.

Интересная проблема — как из поста вытащить контакт заказчика? В Telegram‑каналах с заказами обычно указывают @username или t.me/username. Но надо отфильтровать юзернейм самого канала и ботов:

def _extract_contact(self, text, channel_username):    usernames = re.findall(r"@([a-zA-Z0-9_]{5,32})", text)    for u in usernames:        if u.lower() != channel_username.lower() and not u.lower().endswith("bot"):            return f"@{u}"    # Fallback: t.me/username    links = re.findall(r"(?:t\.me|telegram\.me)/([a-zA-Z0-9_]{5,32})", text)    ...

Если контакт не найден — пост тихо пропускается. Нет способа связаться с заказчиком — нет смысла показывать такой заказ.

Так выглядят заказы в единой ленте
Результат работы скраперов: заказы с Kwork и Freelance в единой ленте

Результат работы скраперов: заказы с Kwork и Freelance в единой ленте

Дедупликация: MD5-хэши и bulk checking

Каждый заказ идентифицируется MD5-хэшем его URL. Почему MD5, а не UUID? Потому что один и тот же URL всегда даёт один и тот же хэш, и это гарантирует идемпотентность: если скрапер трижды увидел один проект, в базе он будет один раз.

def get_url_hash(url: str) -> str:    return hashlib.md5(url.encode("utf-8")).hexdigest()

Каждые 15 секунд (интервал APScheduler) скрапер получает из БД множество всех существующих ID одним SQL‑запросом. На стороне скрапера if url_hash in existing_ids: continue мгновенно отсекает дубликаты ещё до обращения к детальным страницам.

А что если скрапер вернул 50 заказов, а из них 45 уже в базе? execute_batch в юзкейсе сначала делает filter_existing_ids — один IN‑запрос вместо 50 отдельных SELECT‑ов. Новые сохраняет и отправляет уведомления, существующие тихо обновляет метаданные (число откликов, просмотров) без повторных алертов.

SSE вместо WebSocket: event bus на asyncio.Condition

Фронтенд на Vue 3 должен узнавать о новых заказах мгновенно. WebSocket для этого избыточен — мне не нужен двусторонний канал, достаточно push от сервера к клиенту.

SSE (Server‑Sent Events) проще: обычный HTTP, автореконнект из коробки в браузере, работает через любой прокси. Осталось решить, как передать сигнал от скрапера (который работает в APScheduler) к SSE‑эндпоинту (который работает в FastAPI). Оба живут в одном процессе.

Redis pub/sub? Можно, но зачем тащить внешнюю зависимость для in‑process коммуникации? Весь event bus уместился в 38 строк:

_condition = asyncio.Condition()_revision = 0async def notify_new_gigs():    global _revision    async with _condition:        _revision += 1        _condition.notify_all()async def wait_until_new(last_seen, timeout=30.0):    async with _condition:        try:            await asyncio.wait_for(                _condition.wait_for(lambda: _revision > last_seen),                timeout=timeout            )            return _revision        except asyncio.TimeoutError:            return None

Скрапер вызывает notify_new_gigs(). Все подключённые SSE‑клиенты спят на wait_until_new() — и мгновенно просыпаются. Монотонный счётчик revision решает проблему потерянных событий: клиент знает свой lastseen и при реконнекте получит сигнал, если ревизия изменилась.

SSE‑эндпоинт со стороны FastAPI:

async def sse_gigs_stream(request: Request) -> StreamingResponse:    async def generate():        yield ": connected\n\n"        last_seen = get_revision()        while True:            if await request.is_disconnected():                break            new_rev = await wait_until_new(last_seen, timeout=30.0)            if new_rev is not None:                last_seen = new_rev                yield f"id: {new_rev}\ndata: {{\"type\": \"new_gigs\"}}\n\n"            else:                yield ": keepalive\n\n"    return StreamingResponse(generate(), media_type="text/event-stream", ...)

Keepalive‑комментарии каждые 30 секунд предотвращают таймаут прокси и браузера. Лимит 5 соединений на IP защищает от случайного исчерпания ресурсов. Заголовок X-Accel-Buffering: no говорит Nginx/Caddy не буферизовать поток.

Если когда‑нибудь приложение разрастётся до нескольких процессов — заменю asyncio.Condition на Redis pub/sub. Но пока один процесс, и внешняя зависимость не нужна.

Rate limiting для Telegram API

Telegram Bot API имеет жёсткие лимиты: не больше ~30 сообщений в секунду глобально, ~1 сообщение в секунду на одного пользователя. Если скрапер притащил 50 новых заказов и 20 пользователей подписаны на эту категорию — это потенциально 1000 сообщений. Без throttling бот получит TelegramRetryAfter.

Решение — per‑user delay с отслеживанием последнего времени отправки:

class TelegramNotificationService:    _last_send_time = {}    async def _wait_rate_limit(self, telegram_id: int):        now = time.time()        last = self._last_send_time.get(telegram_id, 0.0)        delay = 1.1 - (now - last)        if delay > 0:            self._last_send_time[telegram_id] = now + delay            await asyncio.sleep(delay)        else:            self._last_send_time[telegram_id] = now

1.1 секунды между сообщениями одному юзеру с небольшим запасом. Если Telegram всё равно вернул TelegramRetryAfter — три попытки с asyncio.sleep(retry_err.retry_after).

Оркестрация: всё в одном asyncio.run

Весь стек запускается одним вызовом asyncio.run(start_services()). В одном event loop крутятся:

  1. FastAPI (через uvicorn.Server, не subprocess)

  2. Aiogram bot polling

  3. Support bot polling

  4. APScheduler с тремя скраперами (каждые 15 секунд)

  5. Telethon userbot

  6. SMTP‑сервер

  7. SSE‑стримы

На Windows для psycopg нужна WindowsSelectorEventLoopPolicy — без неё asyncio крашится при работе с сокетами. Такая строчка в начале main.py сэкономила мне часов пять дебага.

Критические ошибки ловятся кастомным TelegramAdminLogHandler и улетают мне в личку в Telegram как logging.CRITICAL. Упал скрапер в три ночи — я узнаю через секунду, а не утром.

Деплой: Docker multi‑stage + Caddy

# Stage 1: Build Vue 3 FrontendFROM node:20-alpine AS frontend-builderWORKDIR /frontendCOPY frontend/package*.json ./RUN npm installCOPY frontend/ ./RUN npm run build# Stage 2: Run Python BackendFROM python:3.13-slimWORKDIR /appCOPY requirements.txt .RUN pip install --no-cache-dir -r requirements.txtCOPY src/ ./srcCOPY --from=frontend-builder /frontend/dist ./static_origin

Фронтенд собирается в первой стадии и копируется как статика. Caddy раздаёт статические файлы напрямую (с Cache-Control: immutable для Vite‑ассетов с хэшами в именах), а всё остальное проксирует на FastAPI.

SMTP‑порт 25 маппится на внутренний 2525: ports: - "25:2525" в docker‑compose. Внешний MX‑запись указывает на сервер, и порт 25 принимает входящую почту от Мейла и Яндекса.

Что я бы сделал иначе

SMTP без TLS — сейчас SMTP‑сервер принимает plaintext. Для внутреннего приёма пересланной почты это работает (MTA отправителя сам устанавливает TLS до моего сервера через STARTTLS на уровне OS/firewall), но для продакшена стоит добавить STARTTLS прямо в хэндлер.

Один процесс — пока всё помещается в одном asyncio loop, но если нагрузка вырастет, скраперы стоит вынести в отдельный воркер с Redis pub/sub вместо asyncio.Condition.

Регулярки в парсерах — хрупкие. Если биржа изменит HTML‑шаблон письма, парсер сломается. Но на практике за полгода шаблоны менялись один раз (Kwork добавил penalty‑блок), и фикс занял 15 минут. Для этого масштаба ML‑классификатор — оверинжиниринг.

Нет graceful shutdown для SMTP — при деплое текущие SMTP‑соединения обрывает Docker. На практике письма от бирж идут по одному, и потеря маловероятна, но asyncio.Event для корректного завершения не помешал бы.


Если у кого‑то похожая задача — собрать несколько внешних источников в единый поток с уведомлениями — надеюсь, описанные решения будут полезны. Особенно SMTP‑сервер на asyncio — штука простая, а покрывает удивительно много юзкейсов, где нужен real‑time приём email без тяжёлой инфраструктуры.

Вопросы, замечания, альтернативные подходы — добро пожаловать в комментарии.

ссылка на оригинал статьи https://habr.com/ru/articles/1064840/