Я фрилансер, и я ненавижу мониторить биржи. 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 строк на весь доменный слой.

@dataclass
class 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 = 0

async 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 Frontend
FROM node:20-alpine AS frontend-builder
WORKDIR /frontend
COPY frontend/package*.json ./
RUN npm install
COPY frontend/ ./
RUN npm run build

# Stage 2: Run Python Backend
FROM python:3.13-slim
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY src/ ./src
COPY --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 без тяжёлой инфраструктуры.

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