"""FastSub — официальный Python SDK.

Установка:
    pip install httpx
    # положите этот файл рядом со своим ботом

Быстрый старт (2 минуты, без единого клика в нашем боте):

    import asyncio
    from fastsub import FastSub

    async def main():
        # 1. Регистрируем бота и получаем ключ. confirm_key берётся один раз
        #    в нашем боте: Интеграция → «Ключ подтверждения».
        async with FastSub.bootstrap(
            confirm_key="fsc_...",
            bot_token="123456:ABC-DEF",
        ) as fs:
            print("ключ:", fs.api_key)

            # 2. Просим задания для пользователя. Всё, что знаете о нём,
            #    передавайте — откроются более дорогие заказы.
            answer = await fs.request_op(
                user_id=123456789,
                has_telegram_premium=True,
            )
            if answer.has_tasks:
                for task in answer.tasks:
                    print(task.title, task.url)
            else:
                print("нет заданий:", answer.explain)

    asyncio.run(main())

Уже есть ключ — просто:

    async with FastSub(api_key="fsp_live_...") as fs:
        ...

Два уровня интеграции, одни и те же методы:

  * «под ключ» (get_links=false в настройках бота) — блок спонсоров отправляем
    мы, вам достаточно request_op() и check_task();
  * «полный» (get_links=true) — вы получаете ссылки и рисуете блок сами.

Ошибки: всё, что вернул сервер не-2xx, поднимается как FastSubError с полем
`detail` (текст для человека) и `status`.
"""

from __future__ import annotations

import asyncio
import contextlib
import logging
import os
import sys
import time
from collections import OrderedDict
from collections.abc import AsyncIterator, Iterator, Sequence
from dataclasses import dataclass, field
from typing import Any

import httpx

__version__ = "2.5.0"

__all__ = [
    "DEFAULT_BASE_URL",
    "EVENT_CLASSES",
    "START_PREFIX",
    "Balance",
    "FastSub",
    "FastSubAdvertiser",
    "FastSubError",
    "FastSubMiddleware",
    "FastSubSync",
    "IssueStatus",
    "OpAnswer",
    "OrderApproved",
    "OrderCompleted",
    "OrderPaused",
    "OrderRejected",
    "OrderResumed",
    "OrderSubscriber",
    "ResourceExpired",
    "ResourceIssued",
    "ResourcePaid",
    "ResourceReverted",
    "ResourceSubscribed",
    "ResourceUnsubscribed",
    "ResourceVerified",
    "SubscriptionCheck",
    "Task",
    "TaskStatus",
    "WebhookEvent",
    "WebhookEvents",
    "fastsub_aiohttp_handler",
    "fastsub_check_router",
    "fastsub_flask_blueprint",
    "fastsub_webhook_router",
    "sponsor_block",
    "start_payload",
    "telebot_gate",
    "verify_webhook",
    "web_step_block",
]

DEFAULT_BASE_URL = "https://fastsub.org"
DEFAULT_TIMEOUT = 20.0

#: Сколько `link_id` принимает `check_many` за один запрос. Больше — это уже
#: выгрузка, а для неё есть статистика.
MAX_CHECK_MANY = 50

#: Свой логгер, ничего никуда не пишущий по умолчанию, — как и положено
#: библиотеке. Включается одной строкой на стороне бота:
#:
#:     logging.getLogger("fastsub").setLevel(logging.DEBUG)
#:
#: Без него «спонсоры не показываются» — это всё, что есть у партнёра: SDK
#: намеренно глотает наши сбои, чтобы не ломать чужого юзера, и до сих пор
#: глотал их молча.
logger = logging.getLogger("fastsub")


def _int_header(headers: Any, name: str) -> int | None:
    """Числовой заголовок или None.

    Заголовка может не быть вовсе, а `Retry-After` вдобавок разрешает вместо
    секунд HTTP-дату. Разбирать дату не будем: лучше честное «не знаю», чем
    число, означающее год.
    """
    if headers is None:
        return None
    try:
        return int(str(headers.get(name)).strip())
    except (AttributeError, TypeError, ValueError):
        return None


class FastSubError(RuntimeError):
    """Сервер ответил ошибкой. `detail` можно показывать пользователю.

    При 429 сервер пишет в заголовки, сколько ждать и какой лимит вы выбрали, —
    и эти числа терялись: до вашего кода доезжало «слишком много запросов» без
    единого намёка, повторять через секунду или через час.

    Своим повторам они и не нужны — `_send` читает `Retry-After` сам, — но
    повторы кончаются. После `max_retries` исключение выходит наружу, и решать,
    когда пробовать снова, приходится уже вам:

        except FastSubError as exc:
            if exc.rate_limited:
                await asyncio.sleep(exc.retry_after or 60)
    """

    def __init__(
        self,
        status: int,
        detail: str,
        payload: dict[str, Any] | None = None,
        *,
        headers: Any = None,
    ) -> None:
        super().__init__(f"[{status}] {detail}")
        self.status = status
        self.detail = detail
        self.payload = payload or {}
        #: Секунды до следующей осмысленной попытки (заголовок `Retry-After`).
        self.retry_after = _int_header(headers, "Retry-After")
        #: Сколько запросов за сколько секунд вам позволено — `X-RateLimit-*`.
        self.rate_limit = _int_header(headers, "X-RateLimit-Limit")
        self.rate_window = _int_header(headers, "X-RateLimit-Window")

    @property
    def rate_limited(self) -> bool:
        """Упёрлись в лимит. В отличие от прочих ошибок, тот же самый запрос
        позже пройдёт — менять в нём нечего, надо подождать `retry_after`."""
        return self.status == 429


@dataclass(slots=True)
class Task:
    """Одно задание для пользователя."""

    link_id: str
    title: str
    button_name: str
    type: str
    task_type: str
    reward_for_publisher: str
    url: str | None = None
    invite_link: str | None = None
    start_link: str | None = None
    retention_bonus_rub: str = "0"
    raw: dict[str, Any] = field(default_factory=dict)

    @property
    def link(self) -> str | None:
        """Единая ссылка: приглашение в канал или запуск бота."""
        return self.url or self.invite_link or self.start_link

    @property
    def total_reward(self) -> str:
        """Выплата вместе с надбавкой за удержание — то, что реально придёт."""
        from decimal import Decimal

        total = Decimal(self.reward_for_publisher) + Decimal(self.retention_bonus_rub)
        return f"{total:.4f}"


@dataclass(slots=True)
class IssueStatus:
    """Состояние одной выдачи — ответ check_resource / элемент check_task.

    Раньше эти методы отдавали голый словарь: чтобы понять, засчитано ли
    задание, приходилось помнить названия полей и держать в голове, какие
    статусы считаются выполненными. Теперь это `done` и `pending`.
    """

    link_id: str
    status: str
    reason: str | None = None
    title: str | None = None
    username: str | None = None
    type: str | None = None
    button_name: str | None = None
    reward_for_publisher: str = "0"
    retention_bonus_rub: str = "0"
    payout_state: str = ""
    hold_until: str | None = None
    is_own: bool = False
    raw: dict[str, Any] = field(default_factory=dict)

    #: Статусы, при которых юзер своё сделал.
    DONE = ("subscribed", "verified", "paid")
    #: …а при этих — уже не сделает: задание закрыто.
    CLOSED = ("expired", "unsubscribed", "reverted", "invalid")

    @property
    def done(self) -> bool:
        """Задание выполнено — можно пускать юзера дальше."""
        return self.status in self.DONE

    @property
    def closed(self) -> bool:
        """Задание закрыто и выполнено уже не будет."""
        return self.status in self.CLOSED

    @property
    def pending(self) -> bool:
        """Ждём действия юзера."""
        return not self.done and not self.closed


@dataclass(slots=True)
class TaskStatus:
    """Ответ check_task: состояние всех выдач одного блока."""

    task_id: str
    resources: list[IssueStatus] = field(default_factory=list)
    raw: dict[str, Any] = field(default_factory=dict)

    @property
    def all_done(self) -> bool:
        """Всё выполнено — открывайте доступ."""
        return bool(self.resources) and all(r.done for r in self.resources)

    @property
    def remaining(self) -> list[IssueStatus]:
        """Что юзеру ещё осталось. Закрытые сюда не попадают: их не сделать."""
        return [r for r in self.resources if r.pending]


@dataclass(slots=True)
class SubscriptionCheck:
    """Ответ check_subscription — живой опрос Telegram."""

    link_id: str
    subscribed: bool
    status: str
    checked_live: bool = False
    reason: str | None = None
    raw: dict[str, Any] = field(default_factory=dict)

    @property
    def retry_is_pointless(self) -> bool:
        """Повторять этот запрос смысла нет — ответ не изменится.

        `invalid_user` — переданного id не существует, проблема в том, что вы
        отправили; `offer_withdrawn` — задание отозвано. Всё остальное имеет
        смысл переспросить позже.
        """
        return self.reason in ("invalid_user", "offer_withdrawn")


#: Что написать под блоком спонсоров, когда юзера дальше не пускаем.
DEFAULT_GATE_MESSAGE = "Выполните задания выше, чтобы продолжить."

#: Префикс наших собственных коллбэков. Их обрабатывает `fastsub_check_router`,
#: и гейт обязан их пропускать: иначе кнопка, которая закрывает задания, сама
#: же и не доезжает до обработчика.
OWN_CALLBACK = "fastsub:"


def _user_signals(user: Any) -> dict[str, Any]:
    """Признаки аудитории, которые уже есть в апдейте.

    Не передать их — значит выбросить единственное, на чём держится языковой
    и premium-таргетинг: заказы с ним просто не покажутся, и выглядеть это будет
    как «дорогих заданий нет».
    """
    return {
        "language_code": getattr(user, "language_code", None),
        "has_telegram_premium": getattr(user, "is_premium", None),
        "has_username": bool(getattr(user, "username", None)) or None,
    }


def _answer_target(event: Any) -> Any:
    """Куда писать: само сообщение или сообщение под нажатой кнопкой."""
    message = getattr(event, "message", None)
    if message is not None and hasattr(message, "answer"):
        return message
    return event if hasattr(event, "answer") else None


async def _show_gate(
    answer: OpAnswer,
    *,
    target: Any,
    gate_message: str,
    send_block: bool,
) -> bool:
    """Показать то, что положено, и сказать, пускать ли дальше.

    Одно место на `op()` и на middleware: два одинаковых гейта разъезжаются на
    первой же правке, а расходятся они молча — один пускает, другой нет.
    """
    if answer.needs_web_step and target is not None:
        text, markup = web_step_block(answer)
        await target.answer(text, reply_markup=markup)
        return False
    if not answer.has_tasks:
        return True
    # В режиме «под ключ» блок уже отправлен нами — второй раз нельзя.
    if answer.delivered:
        return False
    if send_block and target is not None:
        text, markup = sponsor_block(answer, gate_message)
        if markup is not None:
            await target.answer(text, reply_markup=markup)
    return False


#: Насколько часто одному и тому же юзеру имеет смысл повторять `request_op`.
#: Ноль в конструкторе отключает.
DEFAULT_COOLDOWN = 2.0

#: Сколько юзеров помним ради этого. Бот с миллионом юзеров не должен из-за
#: защиты от спама съесть память под них всех.
COOLDOWN_CAPACITY = 4096


class _RecentAnswers:
    """Кто спрашивал только что — и что мы ему ответили.

    Зачем это в клиенте. В режиме «под ключ» каждый `request_op` — это ещё и
    сообщение юзеру от бота партнёра. Один хендлер, повешенный на все апдейты,
    или обычный цикл повторов превращаются в четыре одинаковых блока подряд, и
    юзер жмёт «пожаловаться» — а прилетает это боту партнёра, не нам. Сервер от
    этого защищён своим сторожем, но там сообщение уже не отправлено и партнёр
    видит только, что «ничего не пришло». Здесь же повтор просто получает тот
    же ответ, что и секунду назад.

    Просроченные записи чистятся по ходу, а не по таймеру: у SDK нет своего
    цикла, и заводить его ради словаря — плохой обмен.
    """

    __slots__ = ("_seconds", "_seen")

    def __init__(self, seconds: float) -> None:
        self._seconds = seconds
        self._seen: OrderedDict[tuple[int, str], tuple[float, OpAnswer]] = OrderedDict()

    def get(self, user_id: int, mode: str) -> OpAnswer | None:
        if self._seconds <= 0 or mode == "fresh":
            # `fresh` — это явная просьба собрать список заново. Отдать на неё
            # прошлый ответ значит не сделать ровно то, о чём попросили.
            return None
        found = self._seen.get((user_id, mode))
        if found is None:
            return None
        at, answer = found
        if time.monotonic() - at > self._seconds:
            self._seen.pop((user_id, mode), None)
            return None
        return answer

    def put(self, user_id: int, mode: str, answer: OpAnswer) -> None:
        if self._seconds <= 0 or mode == "fresh":
            return
        self._seen[(user_id, mode)] = (time.monotonic(), answer)
        self._seen.move_to_end((user_id, mode))
        while len(self._seen) > COOLDOWN_CAPACITY:
            self._seen.popitem(last=False)


@dataclass(slots=True)
class Balance:
    """Деньги на счёте и окно проверки отписок.

    Начисление за подписку приходит на баланс сразу и выводится без ожидания:
    дата сброса ниже — это не «когда придут деньги», а конец недели, до которого
    отписка юзера списывает начисление обратно. Вопрос «когда придут деньги»
    имеет один ответ: уже на балансе.
    """

    balance_rub: str
    hold_rub: str
    payout_hold_rub: str
    total_earned_rub: str
    #: Задолженность: начисления за отписавшихся, которые вы уже вывели.
    #: Гасится из следующих подписок автоматически.
    debt_rub: str = "0"
    #: ISO-время ближайшего сброса холда. None — холд на платформе выключен.
    hold_until: str | None = None
    hold_period: str = "week"
    hold_note: str = ""
    raw: dict[str, Any] = field(default_factory=dict)

    @property
    def available_rub(self) -> str:
        """Сколько можно вывести прямо сейчас."""
        return self.balance_rub


@dataclass(slots=True)
class OpAnswer:
    """Ответ на request_op."""

    ok: bool
    tasks: list[Task]
    reason: str | None = None
    delivered: bool = False
    onboarding_url: str | None = None
    onboarding_kind: str | None = None
    #: Чем открывать веб-шаг: "url" (кнопка-ссылка) или "miniapp" (web_app).
    onboarding_open: str | None = None
    #: Чем открывать кнопки заданий. "miniapp" приходит незнакомому юзеру, чьи
    #: кнопки завёрнуты в наш редирект: страница снимет гео, попросит Telegram
    #: открыть канал и закроется. Прямые t.me так открывать нельзя.
    links_open: str | None = None
    availability: dict[str, Any] = field(default_factory=dict)
    raw: dict[str, Any] = field(default_factory=dict)

    @property
    def has_tasks(self) -> bool:
        return bool(self.tasks)

    @property
    def needs_web_step(self) -> bool:
        """Нужно отправить юзера на страницу (анкета или умный редирект)."""
        return bool(self.onboarding_url)

    @property
    def explain(self) -> str:
        """Человеческое объяснение, почему заданий мало или нет."""
        text = self.availability.get("explain")
        if text:
            return str(text)
        return {
            "no_tasks": "Подходящих заданий сейчас нет.",
            "bot_not_moderated": "Бот ещё не прошёл модерацию.",
            "onboarding_required": "Пользователь ещё не прошёл веб-шаг.",
        }.get(self.reason or "", self.reason or "")


@dataclass(slots=True)
class WebhookEvent:
    """Событие, пришедшее на ваш webhook.

    Весь остальной SDK отдаёт разобранные объекты, а событие приходило сырым
    словарём — и `event["payout_reversed"]` требовало помнить, что `true` здесь
    значит «выплату **забрали**», а не «выплата прошла». Перепутать знак в этом
    месте — значит начислить там, где надо списать.

    Поэтому три вопроса, ради которых обработчик и пишут, вынесены в свойства:
    `in_hold`, `money_is_mine`, `money_was_taken`.

    Словарём событие быть не перестало: `event["event"]`, `event.get(...)` и
    `in` работают как раньше, так что обработчики, написанные до появления
    этого класса, менять не нужно.
    """

    event: str
    link_id: str
    status: str
    user_id: int | None = None
    task_id: str | None = None
    publisher_payout_rub: str = "0"
    payout_state: str = ""
    payout_reversed: bool = False
    hold_until: str | None = None
    verified_at: str | None = None
    subscribed_at: str | None = None
    unsubscribed_at: str | None = None
    timestamp: str | None = None
    #: `true` у события, посланного через `webhook_test()`. Начислять по нему
    #: ничего нельзя: за ним нет ни юзера, ни денег.
    test: bool = False
    #: `X-FastSub-Delivery-Id` — единственное, чем различаются повтор и новое
    #: событие. Приёмники SDK дедуплицируют по нему сами.
    delivery_id: str | None = None
    raw: dict[str, Any] = field(default_factory=dict)

    # ---------- чьи сейчас деньги ----------

    @property
    def in_hold(self) -> bool:
        """Начислено, но до `hold_until` могут забрать: юзер ещё может отписаться."""
        return self.payout_state == "held"

    @property
    def money_is_mine(self) -> bool:
        """Холд дожил до конца — деньги ваши, отписка их уже не заберёт."""
        # `resource.verified` — это и есть «холд прошёл»; состояние проверяем
        # тоже, потому что то же самое верно для отписки после конца холда.
        return self.event == "resource.verified" or self.payout_state == "released"

    @property
    def money_was_taken(self) -> bool:
        """Выплату отозвали: отписка внутри холда или отмена начисления."""
        return self.payout_reversed

    # ---------- совместимость со словарём ----------

    def __getitem__(self, key: str) -> Any:
        return self.raw[key]

    def get(self, key: str, default: Any = None) -> Any:
        return self.raw.get(key, default)

    def __contains__(self, key: object) -> bool:
        return key in self.raw


def _retry_worthy(status: int) -> bool:
    """Повторять стоит ровно три случая и ни одного больше.

      * `429` — мы сами попросили подождать, и в заголовке написано сколько;
      * `5xx` — на нашей стороне, у вас всё правильно;
      * обрыв соединения — сеть (ловится отдельно, это не статус).

    `4xx` не повторяется никогда: запрос не станет правильнее оттого, что вы
    отправите его ещё раз, а `/request-op` при повторе внутри TTL вернёт тот же
    список — то есть тихо съест попытку.
    """
    return status == 429 or status >= 500


def _next_delay(resp: httpx.Response | None, delay: float) -> float:
    """Сколько ждать перед следующей попыткой: 429 говорит это сам."""
    if resp is not None and resp.status_code == 429:
        return float(resp.headers.get("Retry-After") or delay)
    return delay


# ---------------------------------------------------------------------------
# Событие своим классом на каждый тип
# ---------------------------------------------------------------------------
#
# Обработчик почти всегда начинается с разбора: что именно произошло. Строкой
# это `if event.event == "resource.verified"` — а опечатка в строке не ошибка,
# она просто ветка, которая никогда не выполнится, и заметить её можно только
# по недосчитанным деньгам.
#
# Классы дают `match` с проверкой на этапе разбора и понятные имена в трассе:
#
#     match event:
#         case ResourceVerified():
#             await credit(event.user_id, event.publisher_payout_rub)
#         case ResourceUnsubscribed() | ResourceReverted() if event.money_was_taken:
#             await debit(event.user_id, event.publisher_payout_rub)
#         case _:
#             log.info("новое событие: %s", event.event)
#
# Всё, что раньше работало со словарём и с базовым `WebhookEvent`, работает и
# сейчас: это тот же класс, только с уточнённым типом.


@dataclass(slots=True)
class ResourceIssued(WebhookEvent):
    """Задание выдано юзеру. Денег пока нет."""


@dataclass(slots=True)
class ResourceSubscribed(WebhookEvent):
    """Юзер подписался. Начислено в холд — до конца холда могут забрать."""


@dataclass(slots=True)
class ResourceVerified(WebhookEvent):
    """Холд прошёл. Деньги ваши, поздняя отписка их уже не заберёт."""


@dataclass(slots=True)
class ResourcePaid(WebhookEvent):
    """Выплачено по заявке на вывод."""


@dataclass(slots=True)
class ResourceUnsubscribed(WebhookEvent):
    """Юзер отписался. Внутри холда это значит, что начисление забрали."""


@dataclass(slots=True)
class ResourceExpired(WebhookEvent):
    """TTL выдачи истёк, юзер так и не подписался."""


@dataclass(slots=True)
class ResourceReverted(WebhookEvent):
    """Начисление отозвано."""


# --- События заказа. Приходят рекламодателю, а не паблишеру: у них разные
# ключи и разные эндпоинты, но одна форма конверта и одна проверка подписи.


class OrderApproved(WebhookEvent):
    """Заказ прошёл модерацию, выдача началась."""


class OrderRejected(WebhookEvent):
    """Заказ отклонён. Причина — в поле `reason`."""


class OrderPaused(WebhookEvent):
    """Выдача остановлена: вами, кончившимся бюджетом или мёртвой целью."""


class OrderResumed(WebhookEvent):
    """Выдача возобновлена."""


class OrderCompleted(WebhookEvent):
    """Заказ набрал своё количество."""


class OrderSubscriber(WebhookEvent):
    """Подписчик по заказу пережил холд и оплачен."""


#: Имя события → класс. Незнакомое имя остаётся базовым `WebhookEvent`:
#: события мы иногда добавляем, и старый SDK не должен от этого падать.
EVENT_CLASSES: dict[str, type[WebhookEvent]] = {
    "resource.issued": ResourceIssued,
    "resource.subscribed": ResourceSubscribed,
    "resource.verified": ResourceVerified,
    "resource.paid": ResourcePaid,
    "resource.unsubscribed": ResourceUnsubscribed,
    "resource.expired": ResourceExpired,
    "resource.reverted": ResourceReverted,
    "order.approved": OrderApproved,
    "order.rejected": OrderRejected,
    "order.paused": OrderPaused,
    "order.resumed": OrderResumed,
    "order.completed": OrderCompleted,
    "order.subscriber": OrderSubscriber,
}


class _AsyncClient:
    """Транспорт, общий для клиента паблишера и клиента рекламодателя.

    Раньше повторы жили только в паблишерском клиенте, а рекламодатель делал
    один голый запрос: 429 на `/orders` прилетал исключением там, где рядом,
    в том же файле, лежал готовый цикл. Разница между ними — префикс пути и
    подпись в User-Agent, и это единственное, что осталось раздельным.
    """

    _prefix = "/api/v1"
    _agent = ""

    def __init__(
        self,
        api_key: str,
        *,
        base_url: str = DEFAULT_BASE_URL,
        timeout: float = DEFAULT_TIMEOUT,
        client: httpx.AsyncClient | None = None,
        max_retries: int = 2,
        request_op_cooldown: float = DEFAULT_COOLDOWN,
    ) -> None:
        self.api_key = api_key
        self.base_url = base_url.rstrip("/")
        # Повторяем только 429, 5xx и обрывы связи — см. `_send`. Ноль отключает.
        self.max_retries = max(0, max_retries)
        self._own_client = client is None
        self._client = client or httpx.AsyncClient(timeout=timeout)
        self._recent = _RecentAnswers(max(0.0, request_op_cooldown))

    # ---------- lifecycle ----------

    async def __aenter__(self) -> Any:
        return self

    async def __aexit__(self, *exc: object) -> None:
        await self.close()

    async def close(self) -> None:
        if self._own_client:
            await self._client.aclose()

    # ---------- transport ----------

    @property
    def _headers(self) -> dict[str, str]:
        return {
            "Authorization": f"Bearer {self.api_key}",
            "Content-Type": "application/json",
            "User-Agent": f"fastsub-python/{__version__}{self._agent}",
        }

    async def _send(
        self,
        method: str,
        path: str,
        *,
        body: dict[str, Any] | None = None,
        params: dict[str, Any] | None = None,
        headers: dict[str, str] | None = None,
    ) -> dict[str, Any]:
        """Один запрос с повторами там, где повтор осмыслен — см. `_retry_worthy`."""
        url = f"{self.base_url}{self._prefix}{path}"
        sending = dict(self._headers)
        if headers:
            sending.update(headers)
        delay = 0.5
        last: Exception | None = None
        for attempt in range(self.max_retries + 1):
            try:
                resp = await self._client.request(
                    method, url, json=body, params=params, headers=sending,
                )
            except httpx.HTTPError as exc:
                last = exc
                if attempt >= self.max_retries:
                    raise FastSubError(0, f"сеть недоступна: {exc}") from exc
                logger.warning("%s %s: %s, повтор через %.1f с", method, path, exc, delay)
            else:
                if not _retry_worthy(resp.status_code):
                    return _unwrap(resp)
                if attempt >= self.max_retries:
                    return _unwrap(resp)  # поднимет FastSubError с телом ответа
                delay = _next_delay(resp, delay)
                logger.warning(
                    "%s %s: http %s, повтор через %.1f с",
                    method, path, resp.status_code, delay,
                )
            await asyncio.sleep(min(delay, 30.0))
            delay *= 2
        raise FastSubError(0, f"не удалось выполнить запрос: {last}")

    async def _post(self, path: str, body: dict[str, Any]) -> dict[str, Any]:
        return await self._send("POST", path, body=body)

    async def _get(
        self, path: str, *, params: dict[str, Any] | None = None,
    ) -> dict[str, Any]:
        return await self._send("GET", path, params=params)


class FastSub(_AsyncClient):
    """Асинхронный клиент. Используйте как контекстный менеджер."""

    async def __aenter__(self) -> FastSub:
        return self

    # ---------- ключ из окружения ----------

    @classmethod
    def from_env(
        cls,
        *,
        key_var: str = "FASTSUB_KEY",
        base_url_var: str = "FASTSUB_BASE_URL",
        **kwargs: Any,
    ) -> FastSub:
        """Клиент из переменных окружения.

            fs = FastSub.from_env()      # FASTSUB_KEY, при желании FASTSUB_BASE_URL

        Первой строкой любого бота стоит `os.environ["FASTSUB_KEY"]` — и там же,
        в примерах из интернета, ключ в конце концов оказывается вписан прямо в
        код. Одна строка вместо трёх убирает повод.

        Ключа нет — падаем сразу и с понятным текстом: `KeyError: 'FASTSUB_KEY'`
        посреди первого запроса объясняет куда меньше.
        """
        key = os.environ.get(key_var)
        if not key:
            raise FastSubError(
                0,
                f"переменная окружения {key_var} не задана — положите в неё "
                f"ключ вида fsp_live_… из мини-аппа «Интеграция»",
            )
        base_url = os.environ.get(base_url_var) or DEFAULT_BASE_URL
        return cls(api_key=key, base_url=base_url, **kwargs)

    # ---------- bootstrap: получить ключ, не открывая бота ----------

    @classmethod
    async def bootstrap(
        cls,
        *,
        confirm_key: str,
        bot_token: str,
        name: str | None = None,
        base_url: str = DEFAULT_BASE_URL,
        timeout: float = DEFAULT_TIMEOUT,
    ) -> FastSub:
        """Регистрирует бота (или находит уже свой) и возвращает готовый клиент.

        Повторный запуск безопасен: если бот уже ваш, вернётся он же со своим
        ключом — поэтому это можно звать из деплой-скрипта.
        """
        payload: dict[str, Any] = {"confirm_key": confirm_key, "bot_token": bot_token}
        if name:
            payload["name"] = name
        async with httpx.AsyncClient(timeout=timeout) as tmp:
            resp = await tmp.post(f"{base_url.rstrip('/')}/api/v1/publisher/bots", json=payload)
        data = _unwrap(resp)
        client = cls(api_key=data["api_key"], base_url=base_url, timeout=timeout)
        client.bot = data.get("bot", {})  # type: ignore[attr-defined]
        return client

    @staticmethod
    async def fetch_key(
        *,
        confirm_key: str,
        bot_token: str,
        base_url: str = DEFAULT_BASE_URL,
        timeout: float = DEFAULT_TIMEOUT,
    ) -> str:
        """API-ключ уже зарегистрированного бота."""
        async with httpx.AsyncClient(timeout=timeout) as tmp:
            resp = await tmp.post(
                f"{base_url.rstrip('/')}/api/v1/publisher/keys",
                json={"confirm_key": confirm_key, "bot_token": bot_token},
            )
        return str(_unwrap(resp)["api_key"])

    # ---------- основные методы ----------

    async def request_op(
        self,
        user_id: int,
        *,
        count: int | None = None,
        has_telegram_premium: bool | None = None,
        has_profile_photo: bool | None = None,
        has_username: bool | None = None,
        has_bio: bool | None = None,
        has_stories: bool | None = None,
        has_gifts: bool | None = None,
        language_code: str | None = None,
        exclude_chat_ids: list[int] | None = None,
        mode: str | None = None,
    ) -> OpAnswer:
        """Задания для пользователя. Главный метод.

        `count` — сколько хотите (не больше настройки бота).

        `mode` — что делать с тем, что юзер уже держит:

          * `keep` (по умолчанию) — тот же список до конца TTL. Юзер, вернувшийся
            в бота, видит то, что ему уже выдали;
          * `fresh` — собрать заново, прошлые невыполненные закрыть;
          * `rotate` — как `keep`, но выполненный список не выдаётся повторно;
          * `top_up` — оставить невыполненные и добить новыми до `count`. То, что
            нужно, когда юзер сделал два задания из трёх: он увидит одно своё и
            два новых, а не те же три.

        Остальное — сигналы об аудитории. Всё, что вы уже знаете о юзере из
        своего бота: передавайте, и откроются заказы с таргетингом по этим
        признакам — они дороже. Не передали — считается «неизвестно», и такие
        заказы этому юзеру не покажутся.
        """
        cached = self._recent.get(user_id, mode or "keep")
        if cached is not None:
            return cached

        body: dict[str, Any] = {"user_id": user_id}
        if count is not None:
            body["count"] = count
        if mode:
            body["mode"] = mode
        # Язык — сигнал таргетинга, а не флаг: передавайте как есть, «pt-BR» тоже.
        if language_code:
            body["language_code"] = language_code
        # Не выдавать эти цели в этом ответе (например, где юзер уже есть).
        if exclude_chat_ids:
            body["exclude_chat_ids"] = list(exclude_chat_ids)
        for name, value in (
            ("has_telegram_premium", has_telegram_premium),
            ("has_profile_photo", has_profile_photo),
            ("has_username", has_username),
            ("has_bio", has_bio),
            ("has_stories", has_stories),
            ("has_gifts", has_gifts),
        ):
            # None means «не сообщали» and must stay absent, not become false —
            # false is a claim about the user, and it filters заказы out.
            if value is not None:
                body[name] = value
        data = await self._post("/request-op", body)
        answer = OpAnswer(
            ok=bool(data.get("ok")),
            tasks=[_task(t) for t in (data.get("tasks") or [])],
            reason=data.get("reason"),
            delivered=bool(data.get("delivered")),
            onboarding_url=data.get("onboarding_url"),
            onboarding_kind=data.get("onboarding_kind"),
            onboarding_open=data.get("onboarding_open"),
            links_open=data.get("links_open"),
            availability=data.get("availability") or {},
            raw=data,
        )
        self._recent.put(user_id, mode or "keep", answer)
        return answer

    async def op(
        self,
        event: Any,
        *,
        count: int | None = None,
        gate_message: str = DEFAULT_GATE_MESSAGE,
        send_block: bool = True,
    ) -> bool:
        """Спросить задания, показать что нужно и сказать, пускать ли дальше.

            @dp.message(CommandStart())
            async def start(message: Message):
                if not await fs.op(message):
                    return
                await message.answer("Доступ открыт")

        Одна строка вместо четырёх: `request_op`, разбор ответа, отправка блока
        или веб-шага и решение. То же самое умеет `FastSubMiddleware`, но её надо
        подключать к диспетчеру — а первая интеграция обычно живёт в одном
        хендлере, и ради неё лезть в middleware никто не хочет.

        Принимает `Message` или `CallbackQuery` — всё, у чего есть `from_user` и
        куда можно ответить. Сигналы об аудитории (язык, Premium, юзернейм)
        берутся из самого апдейта: они там уже есть, и не передать их значит
        отказаться от заказов с таргетингом.

        `True` — заданий нет, продолжайте. `False` — юзеру уже показан блок или
        веб-шаг, дальше его пускать не надо.

        Наш сбой возвращает `True`: реклама не должна ломать чужой продукт.
        Для синхронных ботов есть `telebot_gate(fs, bot)` — там нужен ещё и bot,
        потому что у telebot ответ отправляет он, а не апдейт.
        """
        user = getattr(event, "from_user", None)
        if user is None:
            return True
        try:
            answer = await self.request_op(
                user_id=user.id, count=count, **_user_signals(user),
            )
        except FastSubError as exc:
            logger.warning("request_op не удался, пускаем юзера дальше: %s", exc)
            return True
        except Exception:
            logger.exception("непредвиденная ошибка FastSub, пускаем юзера дальше")
            return True

        try:
            return await _show_gate(
                answer,
                target=_answer_target(event),
                gate_message=gate_message,
                send_block=send_block,
            )
        except Exception:
            # Не смогли показать блок — это наша беда, а не юзера.
            logger.exception("не удалось показать блок, пускаем юзера дальше")
            return True

    async def check_task(self, task_id: str) -> TaskStatus:
        """Состояние всего блока разом. Дёшево — читает сохранённые статусы.

            status = await fs.check_task(answer.task_id)
            if status.all_done:
                ...  # пускаем
        """
        return _task_status(await self._post("/check-task", {"task_id": task_id}))

    async def check_subscription(self, link_id: str) -> SubscriptionCheck:
        """Спросить Telegram прямо сейчас. Дороже, но ответ свежий."""
        return _subscription(
            await self._post("/check-subscription", {"link_id": link_id})
        )

    async def check_resource(self, link_id: str) -> IssueStatus:
        """Статус одной выдачи без обращения к Telegram."""
        return _issue_status(await self._post("/check-resource", {"link_id": link_id}))

    async def check_many(self, link_ids: Sequence[str]) -> list[IssueStatus]:
        """Статусы нескольких выдач одним запросом. До 50 за раз.

        `check_task` умеет только весь блок целиком и только по `task_id`.
        Всё остальное — доводка вчерашних выдач, восстановление состояния после
        перезапуска, проверка того, что юзер насобирал за несколько заходов, —
        приходилось делать циклом по одному запросу на `link_id`. Пятьдесят
        выдач это пятьдесят HTTP-запросов и пятьдесят единиц лимита из трёхсот;
        здесь — один и один.

            for state in await fs.check_many(pending_link_ids):
                if state.done:
                    await grant(state.link_id)

        Порядок ответа совпадает с порядком запроса. Неизвестный или чужой
        `link_id` не роняет пакет: у него `status` пустой, а `reason` —
        `not_found` либо `forbidden`.
        """
        ids = list(link_ids)
        if not ids:
            return []
        if len(ids) > MAX_CHECK_MANY:
            raise ValueError(
                f"за раз можно проверить не больше {MAX_CHECK_MANY} link_id, "
                f"а пришло {len(ids)}: нарежьте список",
            )
        data = await self._post("/check-resources", {"link_ids": ids})
        return [_issue_status(item) for item in (data.get("resources") or [])]

    async def me(self) -> dict[str, Any]:
        """Баланс, окно проверки отписок, рейтинг и лимиты вашего ключа."""
        return await self._get("/me")

    async def balance(self) -> Balance:
        """Деньги и конец окна проверки отписок — без разбора всего `/me`.

        Начисление за подписку приходит на баланс сразу и выводится сразу.
        Дата сброса здесь — не «когда придут деньги», а момент, до которого
        отписка юзера списывает начисление обратно.
        """
        return _balance(await self._get("/me"))

    async def hold(self) -> dict[str, Any]:
        """Что под риском отката в ближайший сброс — с разбивкой по границам.

        `balance()` отвечает «сколько можно вывести». Этот запрос — «сколько из
        заработанного ещё могут забрать отписки»: обычно граница одна, но не всегда.
        """
        return await self._get("/hold")

    async def settings(self) -> dict[str, Any]:
        """Настройки бота: спонсоров за запрос, время сброса, фильтры, блок ОП."""
        return await self._get("/settings")

    async def update_settings(self, **changes: Any) -> dict[str, Any]:
        """Изменить настройки бота. Меняется только присланное.

            await fs.update_settings(sponsors_count=3, min_reward_rub="1.5")
            await fs.update_settings(op_block={"text": "Подпишитесь на {count}"})

        Возвращается карточка целиком — её можно сохранить как новое состояние,
        не склеивая с прежним.
        """
        if not changes:
            raise ValueError("нечего менять: передайте хотя бы одну настройку")
        return await self._send("PATCH", "/settings", body=changes)

    async def user_history(self, user_id: int, *, limit: int = 20) -> dict[str, Any]:
        """Что этот пользователь уже выполнял через вас. Одна страница."""
        return await self._get(f"/users/{user_id}/history", params={"limit": limit})

    def iter_user_history(
        self, user_id: int, *, page_size: int = 100,
    ) -> AsyncIterator[dict[str, Any]]:
        """Вся история юзера, страница за страницей.

            async for row in fs.iter_user_history(123456789):
                ...

        Страницы SDK берёт сам: цикл со счётчиком `offset`, который иначе
        пишется в каждом втором интеграционном скрипте, и в каждом втором —
        с ошибкой на последней странице.
        """
        return _iter_pages(
            self._get, f"/users/{user_id}/history", {}, page_size=page_size,
        )

    async def stats(
        self, *, date_from: str | None = None, date_to: str | None = None,
    ) -> dict[str, Any]:
        """Счётчики по статусам и заработок за период (даты «ГГГГ-ММ-ДД»)."""
        params: dict[str, Any] = {}
        if date_from:
            params["from"] = date_from
        if date_to:
            params["to"] = date_to
        return await self._get("/stats", params=params)

    # ---------- contests ----------

    async def contests(
        self, *, status: str | None = None, limit: int = 20, offset: int = 0,
    ) -> dict[str, Any]:
        """Ваши конкурсы, новые сверху: `{"total", "items"}`."""
        return await self._get("/contests", params=_contest_query(status, limit, offset))

    async def contest(self, contest_id: int) -> dict[str, Any]:
        """Карточка конкурса: условия, ссылки на пост и итоги, код жеребьёвки."""
        return await self._get(f"/contests/{contest_id}")

    async def create_contest(self, **payload: Any) -> dict[str, Any]:
        """Создать конкурс — пост в вашем канале с кнопкой «Участвовать».

            await fs.create_contest(
                channel="@my_channel", text="<b>Розыгрыш</b>",
                media_url="https://example.com/prize.jpg", media_type="photo",
                finish_at=datetime(2026, 9, 20, 18, tzinfo=UTC), winners_count=3,
            )

        Наш бот должен быть админом канала с правом публикации, а вы — админом
        канала. По умолчанию конкурс сразу уходит на проверку; `submit=False` —
        оставить черновиком.
        """
        if not payload.get("channel"):
            raise ValueError("нужен channel: @username канала, ссылка t.me/… или id")
        return await self._post("/contests", _contest_body(payload))

    async def update_contest(self, contest_id: int, **changes: Any) -> dict[str, Any]:
        """Изменить конкурс. Меняется только присланное; `None` — «убрать»."""
        if not changes:
            raise ValueError("нечего менять: передайте хотя бы одно поле")
        return await self._send("PATCH", f"/contests/{contest_id}", body=_contest_body(changes))

    async def submit_contest(self, contest_id: int) -> dict[str, Any]:
        """Отправить черновик на проверку — или сразу в расписание."""
        return await self._post(f"/contests/{contest_id}/submit", {})

    async def preview_contest(self, contest_id: int) -> dict[str, Any]:
        """Прислать пост с кнопкой себе в личку — так он выйдет в канале."""
        return await self._post(f"/contests/{contest_id}/preview", {})

    async def finish_contest(self, contest_id: int) -> dict[str, Any]:
        """Подвести итоги сейчас. Отменить нельзя."""
        return await self._post(f"/contests/{contest_id}/finish", {})

    async def cancel_contest(self, contest_id: int) -> dict[str, Any]:
        return await self._post(f"/contests/{contest_id}/cancel", {})

    async def contest_participants(
        self, contest_id: int, *, status: str | None = None, limit: int = 50, offset: int = 0,
    ) -> dict[str, Any]:
        """Участники: номер, статус, победитель ли и какое место."""
        return await self._get(
            f"/contests/{contest_id}/participants",
            params=_contest_query(status, limit, offset),
        )

    async def contest_stats(self, contest_id: int) -> dict[str, Any]:
        """Воронка, задания, заработок и динамика конкурса."""
        return await self._get(f"/contests/{contest_id}/stats")

    async def contest_sponsors(self, contest_id: int) -> dict[str, Any]:
        """Свои задания конкурса: ваши каналы — вступить или забустить."""
        return await self._get(f"/contests/{contest_id}/sponsors")

    async def add_contest_sponsor(
        self, contest_id: int, link: str, *, task_type: str = "subscribe",
    ) -> dict[str, Any]:
        """Добавить своё задание: `task_type` — "subscribe" или "boost"."""
        return await self._post(
            f"/contests/{contest_id}/sponsors", {"link": link, "task_type": task_type},
        )

    async def remove_contest_sponsor(self, contest_id: int, sponsor_id: int) -> dict[str, Any]:
        return await self._send("DELETE", f"/contests/{contest_id}/sponsors/{sponsor_id}")

    # ---------- webhook ----------

    async def configure_webhook(
        self,
        url: str,
        *,
        events: list[str] | None = None,
        rotate_secret: bool = False,
    ) -> dict[str, Any]:
        """Куда слать события. Секрет возвращается один раз — сохраните его.

        Без webhook вы узнаёте о подтверждении подписки только опросом
        `check_task`; с ним — в момент, когда это произошло.
        """
        body: dict[str, Any] = {"url": url, "rotate_secret": rotate_secret}
        if events:
            body["events"] = list(events)
        return await self._post("/webhook/configure", body)

    async def webhook_info(self) -> dict[str, Any]:
        """Текущая настройка webhook (без секрета) и статистика доставок."""
        return await self._get("/webhook")

    async def webhook_test(self) -> dict[str, Any]:
        """Прислать себе фальшивое событие с полем `test: true`.

        Проверять приёмник на настоящей подписке — значит проверять его на живых
        деньгах.
        """
        return await self._post("/webhook/test", {})

    async def on_access_granted(self, callback: Any) -> None:
        """Вызывается роутером, когда юзер выполнил весь блок.

        По умолчанию — короткое сообщение. Замените на своё: `fs.on_access_granted
        = my_coro`, где my_coro принимает aiogram-колбэк.
        """
        message = getattr(callback, "message", None)
        if message is not None:
            await message.answer("Доступ открыт. Продолжайте.")

    async def check_op(
        self, user_id: int, *, callback_query_id: str | None = None
    ) -> dict[str, Any]:
        """Нажали «Проверить» в блоке, который отправили мы (режим «под ключ»).

        Один вызов: перепроверяем подписки, перерисовываем блок, закрываем
        «часики» на кнопке и отвечаем — пускать или нет.

        Возвращает `allowed`: True — открывайте доступ.
        """
        body: dict[str, Any] = {"user_id": user_id}
        if callback_query_id:
            body["callback_query_id"] = callback_query_id
        return await self._post("/check-op", body)

    # ---------- удобные обёртки ----------

    async def wait_done(
        self,
        link_id: str,
        *,
        attempts: int = 6,
        delay: float = 2.0,
        live: bool = True,
    ) -> bool:
        """Ждать подтверждения задания, опрашивая статус.

        Для кнопки «Проверить» это не нужно — там достаточно одного
        check_subscription. Полезно, когда вы ведёте юзера сами.

        `live=False` — читать сохранённый статус (дёшево), `True` — спрашивать
        Telegram каждый раз.
        """
        for attempt in range(attempts):
            if live:
                answer = await self.check_subscription(link_id)
                if answer.subscribed:
                    return True
                # «Такого юзера нет» и «задание отозвано» — ответ окончательный,
                # ждать нечего. Раньше цикл честно выкручивал все попытки.
                if answer.retry_is_pointless:
                    return False
            else:
                # Здесь именно check_resource: он принимает link_id, а check_task
                # ждёт task_id и на link_id отвечал бы «не найдено» все попытки
                # подряд — тихо, всегда False.
                state = await self.check_resource(link_id)
                if state.done:
                    return True
                if state.closed:
                    return False
            if attempt + 1 < attempts:
                await asyncio.sleep(delay)
        return False

async def _iter_pages(
    fetch: Any,
    path: str,
    params: dict[str, Any],
    *,
    page_size: int,
) -> AsyncIterator[dict[str, Any]]:
    """Обойти постраничную ручку целиком, отдавая по одной записи.

    Пагинация у нас в двух видах — `total` у заказов и подписчиков, `has_more`
    у истории юзера, — и партнёр писал цикл со счётчиком под каждую. Здесь обе
    формы разобраны один раз: цикл заканчивается, когда страница пришла пустой,
    когда `has_more` сказал «всё» или когда набрали `total`.

    Пустая страница — обязательное условие выхода, а не подстраховка: сервер,
    который проигнорировал бы `offset`, иначе крутил бы этот цикл вечно.
    """
    size = max(1, min(page_size, 100))
    offset = 0
    while True:
        page = await fetch(path, params={**params, "limit": size, "offset": offset})
        items = page.get("items") or []
        for item in items:
            yield item
        offset += len(items)
        if not items:
            return
        if "has_more" in page and not page["has_more"]:
            return
        total = page.get("total")
        if total is not None and offset >= int(total):
            return


def verify_webhook(secret: str, body: bytes, signature: str) -> bool:
    """Подпись пришедшего webhook — та ли она.

    Без этой проверки ваш эндпоинт принимает «подтверждение подписки» от кого
    угодно, кто знает его адрес. Считается `HMAC-SHA256(secret, тело как есть)`,
    сравнивается со значением заголовка `X-FastSub-Signature`.

    Тело берите сырым: `await request.body()`, а не перекодированный JSON —
    другой порядок ключей или пробелы дадут другую подпись.

        import fastsub

        @app.post("/fastsub/webhook")
        async def hook(request: Request):
            body = await request.body()
            sig = request.headers.get("X-FastSub-Signature", "")
            if not fastsub.verify_webhook(SECRET, body, sig):
                raise HTTPException(401)
            event = json.loads(body)
            ...
    """
    import hashlib
    import hmac

    expected = hmac.new(secret.encode(), body, hashlib.sha256).hexdigest()
    # Сравнение постоянного времени: обычное `==` выдаёт длину общего префикса.
    return hmac.compare_digest(expected, (signature or "").strip())


def _issue_status(raw: dict[str, Any]) -> IssueStatus:
    return IssueStatus(
        link_id=raw.get("link_id", ""),
        status=raw.get("status") or "",
        reason=raw.get("reason"),
        title=raw.get("title"),
        username=raw.get("username"),
        type=raw.get("type"),
        button_name=raw.get("button_name"),
        reward_for_publisher=str(raw.get("reward_for_publisher") or "0"),
        retention_bonus_rub=str(raw.get("retention_bonus_rub") or "0"),
        payout_state=raw.get("payout_state") or "",
        hold_until=raw.get("hold_until"),
        is_own=bool(raw.get("is_own")),
        raw=raw,
    )


def _task_status(raw: dict[str, Any]) -> TaskStatus:
    return TaskStatus(
        task_id=raw.get("task_id") or "",
        resources=[_issue_status(r) for r in (raw.get("resources") or [])],
        raw=raw,
    )


def _balance(raw: dict[str, Any]) -> Balance:
    hold = raw.get("hold") or {}
    return Balance(
        balance_rub=str(raw.get("balance_rub", "0")),
        hold_rub=str(raw.get("hold_rub", "0")),
        payout_hold_rub=str(raw.get("payout_hold_rub", "0")),
        total_earned_rub=str(raw.get("total_earned_rub", "0")),
        debt_rub=str(raw.get("debt_rub", "0")),
        hold_until=hold.get("next_release_at"),
        hold_period=str(hold.get("period") or "week"),
        hold_note=str(hold.get("description") or ""),
        raw=raw,
    )


def _subscription(raw: dict[str, Any]) -> SubscriptionCheck:
    return SubscriptionCheck(
        link_id=raw.get("link_id", ""),
        subscribed=bool(raw.get("subscribed")),
        status=raw.get("status") or "",
        checked_live=bool(raw.get("checked_live")),
        reason=raw.get("reason"),
        raw=raw,
    )


def _task(raw: dict[str, Any]) -> Task:
    return Task(
        link_id=raw.get("link_id", ""),
        title=raw.get("title") or "",
        button_name=raw.get("button_name") or "Подписаться",
        type=raw.get("type") or "channel",
        task_type=raw.get("task_type") or "subscribe",
        reward_for_publisher=str(raw.get("reward_for_publisher") or "0"),
        retention_bonus_rub=str(raw.get("retention_bonus_rub") or "0"),
        url=raw.get("url"),
        invite_link=raw.get("invite_link"),
        start_link=raw.get("start_link"),
        raw=raw,
    )


def _unwrap(resp: httpx.Response) -> dict[str, Any]:
    try:
        data = resp.json()
    except ValueError:
        raise FastSubError(
            resp.status_code, resp.text[:300] or "пустой ответ", headers=resp.headers,
        ) from None
    if resp.status_code >= 400:
        detail = data.get("detail") or data.get("message") or "ошибка запроса"
        if isinstance(detail, list):  # pydantic validation errors
            detail = "; ".join(str(d.get("msg", d)) for d in detail)
        raise FastSubError(resp.status_code, str(detail), data, headers=resp.headers)
    return data if isinstance(data, dict) else {"data": data}


# ==============================================================================
# aiogram middleware — то, из-за чего SDK перестаёт быть «ещё одним клиентом»
# ==============================================================================
#
# Без него в каждом хендлере приходится: позвать request_op, вспомнить про
# веб-шаг, вспомнить что в режиме «под ключ» блок уже отправлен, и решить —
# пускать юзера дальше или нет. Это одни и те же пятнадцать строк, и та копия,
# где их написали неправильно, пускает всех бесплатно.
#
# Подключение (две строки):
#
#     from fastsub import FastSub, FastSubMiddleware
#
#     fs = FastSub(api_key="fsp_live_...")
#     dp.message.middleware(FastSubMiddleware(fs))
#     dp.callback_query.middleware(FastSubMiddleware(fs))
#
# Дальше в любом хендлере:
#
#     @dp.message(CommandStart())
#     async def start(message: Message, fastsub: OpAnswer) -> None:
#         # fastsub уже готов: задания запрошены, веб-шаг учтён
#         if fastsub.has_tasks:
#             ...
#
# Закрыть доступ до выполнения заданий — флагом на хендлере:
#
#     @dp.message(Command("secret"), flags={"fastsub": "gate"})
#     async def secret(message: Message) -> None:
#         # сюда попадём только если заданий нет или они выполнены
#         await message.answer("Секрет")
#
# Что делает флаг:
#   "gate"  — не пускать, пока есть невыполненные задания (блок уже отправлен);
#   "skip"  — не звать API вообще (для /help, платежей, поддержки);
#   по умолчанию — позвать и положить ответ в data, но пропустить дальше.


def web_step_block(answer: OpAnswer) -> tuple[str, Any]:
    """Текст и кнопка для веб-шага. Кнопка, а не ссылка в тексте.

    Ссылка текстом — это «нажми на синее и вернись сам»: её не видно среди
    сообщения, её нельзя открыть мини-аппом, и половина юзеров на ней
    отваливается. Кнопка называет действие и открывается тем способом, который
    выбрал паблишер: `miniapp` — внутри Telegram, выглядит частью его бота;
    `url` — обычной ссылкой в браузере.

    Возвращает `(text, reply_markup)`; markup — `None`, если aiogram не
    установлен, так что текст в любом случае уйдёт.
    """
    quiz = answer.onboarding_kind == "quiz"
    text = (
        "Заполните короткую анкету — это полминуты, и мы вернём вас обратно."
        if quiz
        else "Один шаг — ничего заполнять не нужно, вернётесь сразу сюда."
    )
    label = "Пройти регистрацию" if quiz else "Продолжить"
    url = str(answer.onboarding_url or "")
    if not url:
        return text, None
    try:
        from aiogram.types import (
            InlineKeyboardButton,
            InlineKeyboardMarkup,
            WebAppInfo,
        )
    except ImportError:  # SDK живёт и без aiogram
        return f"{text}\n{url}", None

    button = (
        InlineKeyboardButton(text=label, web_app=WebAppInfo(url=url))
        if answer.onboarding_open == "miniapp"
        else InlineKeyboardButton(text=label, url=url)
    )
    return text, InlineKeyboardMarkup(inline_keyboard=[[button]])


def _task_button(task: Task, *, miniapp: bool) -> Any:
    """Кнопка одного спонсора.

    Мини-аппом открывается только наша страница перехода — Telegram присылает
    `links_open="miniapp"` ровно тогда, когда ссылка ведёт на неё. Обычную t.me
    так открывать нельзя: вебвью её не отдаёт приложению, и юзер упирается в
    пустой экран.
    """
    from aiogram.types import InlineKeyboardButton, WebAppInfo

    link = task.link or ""
    if miniapp:
        return InlineKeyboardButton(
            text=task.button_name, web_app=WebAppInfo(url=link),
        )
    return InlineKeyboardButton(text=task.button_name, url=link)


def sponsor_block(
    answer: OpAnswer,
    gate_message: str = "Выполните задания, чтобы продолжить.",
    *,
    check_button: str | None = "Я подписался",
    task_id: str | None = None,
) -> tuple[str, Any]:
    """Текст и клавиатура блока спонсоров.

    Одна кнопка проверки на весь список, а не по одной на каждого спонсора.
    Раньше было по одной: три спонсора — три «Проверить», и человек жал их по
    очереди, гадая, какая к какому относится. Проверять по одному незачем — мы
    и так спрашиваем про все разом.

    `check_button=None` убирает кнопку совсем: у кого проверка идёт сама
    (webhook или свой опрос), лишняя кнопка только приглашает нажать её раньше
    времени.

    `task_id` кладётся в callback_data, чтобы обработчик проверил весь блок
    одним вызовом. Если его нет, обработчик проверит выданные ссылки по одной.
    """
    try:
        from aiogram.types import InlineKeyboardButton, InlineKeyboardMarkup
    except ImportError:  # SDK живёт и без aiogram
        return gate_message, None

    lines = [gate_message, ""]
    rows = []
    miniapp = answer.links_open == "miniapp"
    for i, task in enumerate(answer.tasks, start=1):
        lines.append(f"{i}. {task.title}")
        if task.link:
            rows.append([_task_button(task, miniapp=miniapp)])
    if check_button:
        target = task_id or (answer.raw or {}).get("task_id") or ""
        rows.append([
            InlineKeyboardButton(
                text=check_button, callback_data=f"fastsub:done:{target}",
            ),
        ])
    return "\n".join(lines), InlineKeyboardMarkup(inline_keyboard=rows)


class FastSubMiddleware:
    """middleware для aiogram: спрашивает задания сама, показывает блок и решает,
    пускать ли юзера дальше.

        dp.message.middleware(FastSubMiddleware(fs))

    По умолчанию она работает как гейт: незнакомого юзера отправляет на веб-шаг,
    юзеру с невыполненными заданиями показывает блок и хендлер не вызывает.
    Именно это имеют в виду, когда вешают мидлварь, — и именно это она раньше
    **не** делала: гейт включался только флагом `@flags(fastsub="gate")` на
    каждом хендлере, а без флага она молча складывала ответ в `data` и
    пропускала дальше. Со стороны это выглядело как «SDK ничего не выдаёт»: ни
    блока, ни анкеты, ни редиректа — при том что API всё это исправно вернул.

    Флаги никуда не делись и по-прежнему сильнее конструктора — ими удобно
    делать исключения:

        @flags(fastsub="skip")      # служебный хендлер: не спрашивать вовсе
        @flags(fastsub="default")   # ответ в data, показываю блок сам

    Ошибки наши и ваши юзера не блокируют: при любом сбое хендлер вызывается как
    обычно — реклама не должна ломать продукт.
    """

    def __init__(
        self,
        client: FastSub,
        *,
        mode: str = "gate",
        count: int | None = None,
        gate_message: str = DEFAULT_GATE_MESSAGE,
        send_block: bool = True,
        key: str = "fastsub",
    ) -> None:
        self.client = client
        #: Что делать, когда у хендлера нет своего флага: `gate` (по умолчанию),
        #: `default` — только положить ответ в data, `skip` — не спрашивать.
        self.mode = mode
        self.count = count
        self.gate_message = gate_message
        # Рисовать блок самим, если API отдал ссылки (режим «полный»). В режиме
        # «под ключ» блок уже отправлен нами, и второй раз его посылать нельзя.
        self.send_block = send_block
        self.key = key

    # Сигнатуру задаёт aiogram: (handler, event, data) без аннотаций.
    async def __call__(self, handler: Any, event: Any, data: dict[str, Any]) -> Any:
        mode = self._mode(data)
        if mode == "skip":
            data[self.key] = None
            return await handler(event, data)

        # Наши собственные кнопки пропускаем не глядя. Мидлварь вешают и на
        # callback_query — так написано и у нас в примерах, — а в режиме гейта
        # она возвращала None на любой апдейт с невыполненными заданиями. То
        # есть нажатие «Проверить» до `fastsub_check_router` не доезжало
        # вообще: кнопка крутила часики и не делала ничего, а заданий,
        # из-за которых гейт закрыт, ровно она и должна была закрыть.
        callback_data = getattr(event, "data", None)
        if isinstance(callback_data, str) and callback_data.startswith(OWN_CALLBACK):
            data[self.key] = None
            return await handler(event, data)

        user = data.get("event_from_user")
        if user is None:
            data[self.key] = None
            return await handler(event, data)

        answer: OpAnswer | None
        try:
            answer = await self.client.request_op(
                user_id=user.id, count=self.count, **_user_signals(user),
            )
        except FastSubError as exc:
            # Наш сбой — не повод останавливать чужой бот.
            logger.warning("request_op не удался, пускаем юзера дальше: %s", exc)
            answer = None
        except Exception:
            # И наша ошибка — тоже не повод. Раньше здесь ловился только
            # FastSubError, а всё остальное — опечатка в нашем коде, неверный
            # объект вместо клиента, сломанный ответ — убивало чужой хендлер.
            # Обещание «реклама не ломает продукт» не должно зависеть от того,
            # какого класса у нас баг.
            logger.exception("непредвиденная ошибка FastSub, пускаем юзера дальше")
            answer = None

        data[self.key] = answer

        if mode != "gate" or answer is None:
            return await handler(event, data)

        # --- gate: дальше только если делать нечего ---
        # Тот же гейт, что у `FastSub.op()`: два одинаковых разъедутся молча —
        # один будет пускать, другой нет.
        try:
            allowed = await _show_gate(
                answer,
                target=self._target(event),
                gate_message=self.gate_message,
                send_block=self.send_block,
            )
        except Exception:
            logger.exception("не удалось показать блок, пускаем юзера дальше")
            allowed = True
        if allowed:
            return await handler(event, data)
        return None

    # ---------- внутреннее ----------

    def _mode(self, data: dict[str, Any]) -> str:
        """Флаг хендлера сильнее режима мидлвари, а тот — сильнее умолчания."""
        flags = data.get("handler")
        flags = getattr(flags, "flags", None) or {}
        value = flags.get("fastsub")
        if value in ("gate", "skip", "default"):
            return str(value)
        return self.mode if self.mode in ("gate", "skip", "default") else "gate"

    @staticmethod
    def _target(event: Any) -> Any:
        """Куда писать: сообщение или сообщение под кнопкой."""
        message = getattr(event, "message", None)
        if message is not None and hasattr(message, "answer"):
            return message
        return event if hasattr(event, "answer") else None


def fastsub_check_router(client: FastSub) -> Any:
    """Готовый роутер для кнопки «Проверить» из блока middleware.

        dp.include_router(fastsub_check_router(fs))

    Обрабатывает только свои коллбэки, поэтому не мешает вашим:

      * `fastsub:done:*` — «Я подписался» под блоком. Одна кнопка на весь
        список: спрашиваем про все задания разом, перерисовываем сообщение,
        оставляя только тех, к кому юзер ещё не подписался, и открываем доступ,
        когда не осталось никого.
      * `fastsub:op:*` — блок «под ключ»: один вызов `/check-op`, всё остальное
        на нашей стороне.
      * `fastsub:check:*` — по одной кнопке на спонсора. Так было до 2.1.2;
        обработчик остался, потому что старые сообщения живут в чатах дальше.

    Когда всё выполнено, вызовется `client.on_access_granted(callback)`; по
    умолчанию он просто пишет юзеру, что доступ открыт — замените на своё:

        fs.on_access_granted = my_handler
    """
    from aiogram import F, Router
    from aiogram.types import CallbackQuery

    router = Router(name="fastsub_check")

    # Режим «под ключ»: блок отправили мы, кнопка наша, и всё за ней делает один
    # вызов. Юзеру ответим мы же — токен, которым отправлено сообщение, ваш, поэтому
    # id колбэка передаём нам.
    @router.callback_query(F.data.startswith("fastsub:op:"))
    async def _check_op(callback: CallbackQuery) -> None:
        if callback.from_user is None:
            return
        try:
            data = await client.check_op(
                callback.from_user.id, callback_query_id=callback.id
            )
        except FastSubError as e:
            await callback.answer(e.detail[:190], show_alert=True)
            return
        if data.get("allowed"):
            # Доступ открыт. Здесь ваш код: показать меню, выдать контент.
            await client.on_access_granted(callback)

    # Одна кнопка на весь блок: спрашиваем про все задания разом, вычёркиваем
    # выполненные и оставляем в сообщении только тех, к кому ещё не подписались.
    @router.callback_query(F.data.startswith("fastsub:done:"))
    async def _check_all(callback: CallbackQuery) -> None:
        if callback.from_user is None:
            return
        task_id = (callback.data or "").split(":")[-1]
        try:
            status = await client.check_task(task_id) if task_id else None
            live = list(status.remaining) if status else []
            if status is None:
                await callback.answer(
                    "Не понимаю, что проверять — попросите список заново.",
                    show_alert=True,
                )
                return
            # `check_task` читает сохранённые статусы: они точны, когда событие
            # о вступлении уже дошло, и отстают, когда нет. Про оставшиеся
            # спрашиваем Telegram живьём — юзер только что нажал кнопку, он ждёт
            # ответа про сейчас, а не про минуту назад.
            still: list[IssueStatus] = []
            for item in live:
                fresh = await client.check_subscription(item.link_id)
                if not fresh.subscribed:
                    still.append(item)
        except FastSubError as e:
            await callback.answer(e.detail[:190], show_alert=True)
            return

        if not still:
            await callback.answer("Готово, спасибо!", show_alert=False)
            await client.on_access_granted(callback)
            return

        await callback.answer(
            f"Осталось: {len(still)}. Подпишитесь и нажмите ещё раз.",
            show_alert=True,
        )
        await _redraw(callback, still)

    async def _redraw(callback: CallbackQuery, remaining: list[Any]) -> None:
        """Перерисовать блок, оставив только невыполненное.

        Список, в котором выполненные строки продолжают висеть, читается как
        «ничего не засчиталось»: человек видит те же три пункта и не понимает,
        что два уже закрыты.
        """
        from aiogram.types import InlineKeyboardButton, InlineKeyboardMarkup

        message = getattr(callback, "message", None)
        if message is None:
            return
        head = (message.text or "").split("\n", 1)[0]
        lines = [head, ""]
        rows = []
        for i, item in enumerate(remaining, start=1):
            title = item.title or item.link_id
            lines.append(f"{i}. {title}")
            link = (item.raw or {}).get("url") or (item.raw or {}).get("invite_link")
            if link:
                rows.append([
                    InlineKeyboardButton(
                        text=item.button_name or "Подписаться", url=link,
                    ),
                ])
        rows.append([
            InlineKeyboardButton(
                text="Я подписался",
                callback_data=(callback.data or "fastsub:done:"),
            ),
        ])
        # Сообщение могли удалить, или оно не изменилось — не наша беда.
        with contextlib.suppress(Exception):
            await message.edit_text(
                "\n".join(lines),
                reply_markup=InlineKeyboardMarkup(inline_keyboard=rows),
            )

    # Старый коллбэк с одной кнопкой на каждого спонсора: блоки, отправленные
    # до обновления, ещё живут в чатах и должны продолжать работать.
    @router.callback_query(F.data.startswith("fastsub:check:"))
    async def _check(callback: CallbackQuery) -> None:
        link_id = (callback.data or "").split(":")[-1]
        try:
            answer = await client.check_subscription(link_id)
        except FastSubError as e:
            await callback.answer(e.detail[:190], show_alert=True)
            return
        if answer.subscribed:
            await callback.answer("Спасибо! Задание зачтено.", show_alert=True)
        else:
            await callback.answer(
                "Пока не видим подписку. Подпишитесь и попробуйте снова.",
                show_alert=True,
            )

    return router


# ---------------------------------------------------------------------------
# Синхронный клиент
# ---------------------------------------------------------------------------
#
# Асинхронный клиент бесполезен половине ботов: pyTelegramBotAPI, telebot и
# ботов на Flask пишут синхронно, и заставлять их ради нас городить asyncio.run
# внутри хендлера — это ровно та причина, по которой SDK не берут.
#
# Класс не дублирует логику: методы и модели те же, отличается только транспорт.
# Что не переносится в синхронный мир — middleware aiogram: он по определению
# асинхронная.


class FastSubSync:
    """То же самое, но без async. Используйте как контекстный менеджер.

        with FastSubSync(api_key="fsp_live_...") as fs:
            answer = fs.request_op(user_id=123456789)
            for task in answer.tasks:
                bot.send_message(chat_id, task.title, ...)
    """

    def __init__(
        self,
        api_key: str,
        *,
        base_url: str = DEFAULT_BASE_URL,
        timeout: float = DEFAULT_TIMEOUT,
        client: httpx.Client | None = None,
        max_retries: int = 2,
        request_op_cooldown: float = DEFAULT_COOLDOWN,
    ) -> None:
        self.api_key = api_key
        self.base_url = base_url.rstrip("/")
        self.max_retries = max(0, max_retries)
        self._own_client = client is None
        self._client = client or httpx.Client(timeout=timeout)
        self._recent = _RecentAnswers(max(0.0, request_op_cooldown))

    @classmethod
    def from_env(
        cls,
        *,
        key_var: str = "FASTSUB_KEY",
        base_url_var: str = "FASTSUB_BASE_URL",
        **kwargs: Any,
    ) -> FastSubSync:
        """Клиент из переменных окружения — см. `FastSub.from_env`."""
        key = os.environ.get(key_var)
        if not key:
            raise FastSubError(
                0,
                f"переменная окружения {key_var} не задана — положите в неё "
                f"ключ вида fsp_live_… из мини-аппа «Интеграция»",
            )
        return cls(
            api_key=key,
            base_url=os.environ.get(base_url_var) or DEFAULT_BASE_URL,
            **kwargs,
        )

    def __enter__(self) -> FastSubSync:
        return self

    def __exit__(self, *exc: object) -> None:
        self.close()

    def close(self) -> None:
        if self._own_client:
            self._client.close()

    # ---------- методы ----------

    def request_op(self, user_id: int, **kwargs: Any) -> OpAnswer:
        """Задания для пользователя. Параметры те же, что у асинхронного."""
        mode = str(kwargs.get("mode") or "keep")
        cached = self._recent.get(user_id, mode)
        if cached is not None:
            return cached
        body = _request_op_body(user_id, kwargs)
        answer = _op_answer(self._send("POST", "/request-op", body=body))
        self._recent.put(user_id, mode, answer)
        return answer

    def check_task(self, task_id: str) -> TaskStatus:
        return _task_status(self._send("POST", "/check-task", body={"task_id": task_id}))

    def check_resource(self, link_id: str) -> IssueStatus:
        return _issue_status(
            self._send("POST", "/check-resource", body={"link_id": link_id})
        )

    def check_subscription(self, link_id: str) -> SubscriptionCheck:
        return _subscription(
            self._send("POST", "/check-subscription", body={"link_id": link_id})
        )

    def check_op(
        self, user_id: int, *, callback_query_id: str | None = None,
    ) -> dict[str, Any]:
        body: dict[str, Any] = {"user_id": user_id}
        if callback_query_id:
            body["callback_query_id"] = callback_query_id
        return self._send("POST", "/check-op", body=body)

    def me(self) -> dict[str, Any]:
        return self._send("GET", "/me")

    def balance(self) -> Balance:
        """Деньги и когда закроется холд — см. `FastSub.balance`."""
        return _balance(self._send("GET", "/me"))

    def op(self, message: Any, bot: Any, *, count: int | None = None) -> bool:
        """Спросить задания, показать блок и сказать, пускать ли дальше.

            @bot.message_handler(commands=["start"])
            def start(message):
                if not fs.op(message, bot):
                    return
                bot.send_message(message.chat.id, "Доступ открыт")

        То же, что `FastSub.op`, но `bot` приходится передавать: у telebot
        отвечает он, а не сам апдейт — у сообщения нет метода, которым можно
        ответить.

        Пишем в тот чат, откуда пришло сообщение, а задания просим для
        `from_user`: в личке это одно и то же число, а в группе — нет, и
        перепутать их значит запросить задания для чата, которого не существует.

        `True` — продолжайте, `False` — блок показан. Наш сбой возвращает
        `True`: реклама не должна ломать чужой продукт.
        """
        user = getattr(message, "from_user", None)
        chat = getattr(message, "chat", None)
        if user is None:
            return True
        allowed = telebot_gate(self, bot, count=count)(
            user.id, getattr(chat, "id", None),
        )
        return bool(allowed)

    def hold(self) -> dict[str, Any]:
        """Что закроется в ближайший сброс — см. `FastSub.hold`."""
        return self._send("GET", "/hold")

    def settings(self) -> dict[str, Any]:
        """Настройки бота — см. `FastSub.settings`."""
        return self._send("GET", "/settings")

    def update_settings(self, **changes: Any) -> dict[str, Any]:
        """Изменить настройки бота — см. `FastSub.update_settings`."""
        if not changes:
            raise ValueError("нечего менять: передайте хотя бы одну настройку")
        return self._send("PATCH", "/settings", body=changes)

    def stats(
        self, *, date_from: str | None = None, date_to: str | None = None,
    ) -> dict[str, Any]:
        params: dict[str, Any] = {}
        if date_from:
            params["from"] = date_from
        if date_to:
            params["to"] = date_to
        return self._send("GET", "/stats", params=params)

    def user_history(self, user_id: int, *, limit: int = 20) -> dict[str, Any]:
        return self._send(
            "GET", f"/users/{user_id}/history", params={"limit": limit},
        )

    # ---------- contests: то же, что у FastSub ----------

    def contests(
        self, *, status: str | None = None, limit: int = 20, offset: int = 0,
    ) -> dict[str, Any]:
        return self._send("GET", "/contests", params=_contest_query(status, limit, offset))

    def contest(self, contest_id: int) -> dict[str, Any]:
        return self._send("GET", f"/contests/{contest_id}")

    def create_contest(self, **payload: Any) -> dict[str, Any]:
        """Создать конкурс — см. `FastSub.create_contest`."""
        if not payload.get("channel"):
            raise ValueError("нужен channel: @username канала, ссылка t.me/… или id")
        return self._send("POST", "/contests", body=_contest_body(payload))

    def update_contest(self, contest_id: int, **changes: Any) -> dict[str, Any]:
        if not changes:
            raise ValueError("нечего менять: передайте хотя бы одно поле")
        return self._send("PATCH", f"/contests/{contest_id}", body=_contest_body(changes))

    def submit_contest(self, contest_id: int) -> dict[str, Any]:
        return self._send("POST", f"/contests/{contest_id}/submit", body={})

    def preview_contest(self, contest_id: int) -> dict[str, Any]:
        return self._send("POST", f"/contests/{contest_id}/preview", body={})

    def finish_contest(self, contest_id: int) -> dict[str, Any]:
        return self._send("POST", f"/contests/{contest_id}/finish", body={})

    def cancel_contest(self, contest_id: int) -> dict[str, Any]:
        return self._send("POST", f"/contests/{contest_id}/cancel", body={})

    def contest_participants(
        self, contest_id: int, *, status: str | None = None, limit: int = 50, offset: int = 0,
    ) -> dict[str, Any]:
        return self._send(
            "GET", f"/contests/{contest_id}/participants",
            params=_contest_query(status, limit, offset),
        )

    def contest_stats(self, contest_id: int) -> dict[str, Any]:
        return self._send("GET", f"/contests/{contest_id}/stats")

    def contest_sponsors(self, contest_id: int) -> dict[str, Any]:
        return self._send("GET", f"/contests/{contest_id}/sponsors")

    def add_contest_sponsor(
        self, contest_id: int, link: str, *, task_type: str = "subscribe",
    ) -> dict[str, Any]:
        return self._send(
            "POST", f"/contests/{contest_id}/sponsors",
            body={"link": link, "task_type": task_type},
        )

    def remove_contest_sponsor(self, contest_id: int, sponsor_id: int) -> dict[str, Any]:
        return self._send("DELETE", f"/contests/{contest_id}/sponsors/{sponsor_id}")

    def configure_webhook(
        self, url: str, *, events: list[str] | None = None, rotate_secret: bool = False,
    ) -> dict[str, Any]:
        body: dict[str, Any] = {"url": url, "rotate_secret": rotate_secret}
        if events:
            body["events"] = list(events)
        return self._send("POST", "/webhook/configure", body=body)

    def check_many(self, link_ids: Sequence[str]) -> list[IssueStatus]:
        """Статусы нескольких выдач одним запросом — см. `FastSub.check_many`."""
        ids = list(link_ids)
        if not ids:
            return []
        if len(ids) > MAX_CHECK_MANY:
            raise ValueError(
                f"за раз можно проверить не больше {MAX_CHECK_MANY} link_id, "
                f"а пришло {len(ids)}: нарежьте список",
            )
        data = self._send("POST", "/check-resources", body={"link_ids": ids})
        return [_issue_status(item) for item in (data.get("resources") or [])]

    def iter_user_history(
        self, user_id: int, *, page_size: int = 100,
    ) -> Iterator[dict[str, Any]]:
        """Вся история юзера страницами — синхронный близнец итератора.

            for row in fs.iter_user_history(123456789):
                ...
        """
        size = max(1, min(page_size, 100))
        offset = 0
        while True:
            page = self._send(
                "GET", f"/users/{user_id}/history",
                params={"limit": size, "offset": offset},
            )
            items = page.get("items") or []
            yield from items
            offset += len(items)
            if not items:
                return
            if "has_more" in page and not page["has_more"]:
                return
            total = page.get("total")
            if total is not None and offset >= int(total):
                return

    def webhook_info(self) -> dict[str, Any]:
        """Текущая настройка webhook (без секрета) и статистика доставок."""
        return self._send("GET", "/webhook")

    def webhook_test(self) -> dict[str, Any]:
        """Прислать себе фальшивое событие с полем `test: true`."""
        return self._send("POST", "/webhook/test", body={})

    def wait_done(
        self,
        link_id: str,
        *,
        attempts: int = 6,
        delay: float = 2.0,
        live: bool = True,
    ) -> bool:
        """Ждать подтверждения задания, опрашивая статус.

        Синхронный близнец `FastSub.wait_done`, вплоть до правил выхода: «такого
        юзера нет» и «задание отозвано» — ответ окончательный, крутить попытки
        не о чем.
        """
        import time

        for attempt in range(attempts):
            if live:
                answer = self.check_subscription(link_id)
                if answer.subscribed:
                    return True
                if answer.retry_is_pointless:
                    return False
            else:
                state = self.check_resource(link_id)
                if state.done:
                    return True
                if state.closed:
                    return False
            if attempt + 1 < attempts:
                time.sleep(delay)
        return False

    # ---------- transport ----------

    def _send(
        self,
        method: str,
        path: str,
        *,
        body: dict[str, Any] | None = None,
        params: dict[str, Any] | None = None,
    ) -> dict[str, Any]:
        """Те же правила повторов, что в асинхронном клиенте, — буквально те же:
        решение «повторять ли» и «через сколько» берётся из общих `_retry_worthy`
        и `_next_delay`, чтобы два транспорта не разъехались в поведении."""
        import time

        url = f"{self.base_url}/api/v1{path}"
        headers = {
            "Authorization": f"Bearer {self.api_key}",
            "Content-Type": "application/json",
            "User-Agent": f"fastsub-python/{__version__} (sync)",
        }
        delay = 0.5
        last: Exception | None = None
        for attempt in range(self.max_retries + 1):
            try:
                resp = self._client.request(
                    method, url, json=body, params=params, headers=headers,
                )
            except httpx.HTTPError as exc:
                last = exc
                if attempt >= self.max_retries:
                    raise FastSubError(0, f"сеть недоступна: {exc}") from exc
                logger.warning("%s %s: %s, повтор через %.1f с", method, path, exc, delay)
            else:
                if not _retry_worthy(resp.status_code):
                    return _unwrap(resp)
                if attempt >= self.max_retries:
                    return _unwrap(resp)
                delay = _next_delay(resp, delay)
                logger.warning(
                    "%s %s: http %s, повтор через %.1f с",
                    method, path, resp.status_code, delay,
                )
            time.sleep(min(delay, 30.0))
            delay *= 2
        raise FastSubError(0, f"не удалось выполнить запрос: {last}")


def telebot_gate(
    fs: FastSubSync,
    bot: Any,
    *,
    count: int | None = None,
) -> Any:
    """Готовая проверка для pyTelegramBotAPI: пускать юзера или показать блок.

    aiogram получает middleware, а синхронные боты — вот это. Возвращает функцию,
    которую зовут в начале хендлера:

        gate = telebot_gate(fs, bot)

        @bot.message_handler(commands=["start"])
        def start(message):
            if not gate(message.from_user.id, message.chat.id):
                return          # блок уже показан, дальше не пускаем
            bot.send_message(message.chat.id, "Доступ открыт")

    Первый аргумент — **юзер**, второй — куда писать. Раньше в примере стоял
    `message.chat.id` и он же уходил в `request_op`: в личке это одно и то же
    число, а в группе — id группы, и задания запрашивались для пользователя,
    которого не существует. Второй аргумент можно не передавать — тогда пишем
    самому юзеру, как и было.

    Наш сбой юзера не блокирует: если API недоступен, функция вернёт True и
    хендлер отработает как обычно. Реклама не должна ломать продукт.
    """
    # Импорт — проверка, что библиотека вообще стоит: без неё понятная
    # ошибка здесь лучше, чем AttributeError на первой кнопке.
    import telebot  # noqa: F401

    def _gate(user_id: int, chat_id: int | None = None) -> bool:
        where = user_id if chat_id is None else chat_id
        try:
            answer = fs.request_op(user_id, count=count)
        except FastSubError as exc:
            logger.warning("request_op не удался, пускаем юзера дальше: %s", exc)
            return True
        except Exception:
            logger.exception("непредвиденная ошибка FastSub, пускаем юзера дальше")
            return True
        if answer.onboarding_url:
            bot.send_message(
                where,
                "Нажмите кнопку ниже, чтобы продолжить.",
                reply_markup=_telebot_link_kb(
                    [("Продолжить", answer.onboarding_url)],
                ),
            )
            return False
        if answer.delivered:
            # Режим «под ключ»: блок отправили мы, рисовать нечего.
            return False
        if not answer.has_tasks:
            return True

        rows = [
            (f"{t.button_name}: {t.title}", link)
            for t in answer.tasks
            if (link := t.link)
        ]
        if not rows:
            return True
        bot.send_message(
            where,
            "Подпишитесь, чтобы продолжить:",
            reply_markup=_telebot_link_kb(rows),
        )
        return False

    return _gate


def _telebot_link_kb(rows: list[tuple[str, str]]) -> Any:
    from telebot import types

    kb = types.InlineKeyboardMarkup()  # type: ignore[no-untyped-call]  # telebot ships no annotations
    for text, url in rows:
        kb.add(types.InlineKeyboardButton(text=text[:60], url=url))
    return kb


def _contest_body(payload: dict[str, Any]) -> dict[str, Any]:
    """Тело конкурса для JSON: даты — ISO-строкой, деньги — строкой.

    Партнёр естественно передаёт `datetime` и `Decimal`, а httpx их не
    сериализует — и падал бы на ровном месте, до сети.
    """
    def plain(value: Any) -> Any:
        if hasattr(value, "isoformat"):
            return value.isoformat()
        if isinstance(value, dict):
            return {k: plain(v) for k, v in value.items()}
        if isinstance(value, (list, tuple)):
            return [plain(v) for v in value]
        if value is None or isinstance(value, (str, int, float, bool)):
            return value
        return str(value)

    return {key: plain(value) for key, value in payload.items()}


def _contest_query(status: str | None, limit: int, offset: int) -> dict[str, Any]:
    params: dict[str, Any] = {"limit": limit, "offset": offset}
    if status:
        params["status"] = status
    return params


def _request_op_body(user_id: int, kwargs: dict[str, Any]) -> dict[str, Any]:
    """Тело /request-op — одно на оба клиента, чтобы они не разъезжались."""
    body: dict[str, Any] = {"user_id": user_id}
    for key in ("count", "mode"):
        if kwargs.get(key) is not None:
            body[key] = kwargs[key]
    if kwargs.get("language_code"):
        body["language_code"] = kwargs["language_code"]
    if kwargs.get("exclude_chat_ids"):
        body["exclude_chat_ids"] = list(kwargs["exclude_chat_ids"])
    for name in (
        "has_telegram_premium", "has_profile_photo", "has_username",
        "has_bio", "has_stories", "has_gifts",
    ):
        value = kwargs.get(name)
        # None — «не сообщали». False — утверждение о юзере, оно отсекает заказы.
        if value is not None:
            body[name] = value
    return body


def _op_answer(data: dict[str, Any]) -> OpAnswer:
    return OpAnswer(
        ok=bool(data.get("ok")),
        tasks=[_task(t) for t in (data.get("tasks") or [])],
        reason=data.get("reason"),
        delivered=bool(data.get("delivered")),
        onboarding_url=data.get("onboarding_url"),
        onboarding_kind=data.get("onboarding_kind"),
        onboarding_open=data.get("onboarding_open"),
        availability=data.get("availability") or {},
        raw=data,
    )


# ---------------------------------------------------------------------------
# Приёмник webhook
# ---------------------------------------------------------------------------
#
# В документации это был пример на двадцать строк, который каждый партнёр
# переписывал у себя — вместе с проверкой подписи, которую половина забывает.
# Готовый роутер убирает и то, и другое.
#
# Доставка у нас «хотя бы один раз», и это не оговорка, а свойство: воркер
# отправляет событие вне транзакции и записывает результат отдельной, поэтому
# упавший между этими шагами отправит его ещё раз. Плюс повтор приходит на
# любой не-2xx — а «обработчик успел начислить и упал на следующей строке»
# выглядит снаружи именно как не-2xx.
#
# Значит, дедуп по `X-FastSub-Delivery-Id` обязателен. Раньше приёмники этот
# заголовок даже не читали и в обработчик не передавали, так что написать дедуп
# самому было нельзя при всём желании: id до партнёра просто не доезжал.


class _SeenDeliveries:
    """Какие доставки уже обработаны — чтобы не начислить бонус дважды.

    Набор в памяти процесса, ограниченный по размеру: переживает обычные
    повторы (они приходят минутами позже), но не переживает рестарт и ничего
    не знает о соседних воркерах. Для одного процесса этого достаточно, для
    нескольких — передайте в приёмник свою функцию, см. `dedupe`.
    """

    __slots__ = ("_capacity", "_seen")

    def __init__(self, capacity: int = 4096) -> None:
        self._capacity = capacity
        self._seen: OrderedDict[str, None] = OrderedDict()

    def __call__(self, delivery_id: str) -> bool:
        """True — доставку видим впервые, её надо обработать."""
        if delivery_id in self._seen:
            self._seen.move_to_end(delivery_id)
            return False
        self._seen[delivery_id] = None
        if len(self._seen) > self._capacity:
            self._seen.popitem(last=False)
        return True


class WebhookEvents:
    """Обработчики по типу события — вместо `if/elif` по строкам.

        events = WebhookEvents()

        @events.on("resource.verified")
        async def paid(event: WebhookEvent) -> None:
            await credit(event.user_id, event.publisher_payout_rub)

        @events.on("resource.unsubscribed", "resource.reverted")
        async def taken_back(event: WebhookEvent) -> None:
            if event.money_was_taken:
                await debit(event.user_id, event.publisher_payout_rub)

        app.include_router(fastsub_webhook_router(SECRET, events))

    Обработчиков на одно событие может быть сколько угодно — позовём все по
    порядку. `on("*")` ловит то, что не разобрали по имени: события мы иногда
    добавляем, и лучше их залогировать, чем потерять молча.
    """

    __slots__ = ("_handlers",)

    def __init__(self) -> None:
        self._handlers: dict[str, list[Any]] = {}

    def on(self, *events: str) -> Any:
        def wrap(fn: Any) -> Any:
            for name in events:
                self._handlers.setdefault(name, []).append(fn)
            return fn

        return wrap

    def handlers_for(self, event: str) -> list[Any]:
        return self._handlers.get(event) or self._handlers.get("*") or []

    async def __call__(self, event: WebhookEvent) -> None:
        for fn in self.handlers_for(event.event):
            result = fn(event)
            if asyncio.iscoroutine(result):
                await result

    def dispatch(self, event: WebhookEvent) -> None:
        """Синхронный вызов — для Flask и прочих синхронных приёмников."""
        for fn in self.handlers_for(event.event):
            result = fn(event)
            if asyncio.iscoroutine(result):
                # Корутину никто не выполнит, а несозданная задача ещё и
                # заругается в лог при сборке мусора. Скажем прямо, в чём дело.
                result.close()
                raise TypeError(
                    f"обработчик {getattr(fn, '__name__', fn)!r} асинхронный, "
                    "а приёмник синхронный — возьмите fastsub_webhook_router "
                    "или сделайте обработчик обычной функцией"
                )


class _BadSignatureError(Exception):
    """Подпись не сошлась: тело подменили или секрет не тот."""


def _dedupe_fn(dedupe: Any) -> Any:
    """`True` — набор в памяти, `False`/`None` — без дедупа, функция — ваша."""
    if dedupe is True:
        return _SeenDeliveries()
    if not dedupe:
        return None
    return dedupe


def _parse_delivery(
    secret: str, body: bytes, headers: Any, dedupe: Any
) -> WebhookEvent | None:
    """Общая часть всех приёмников: подпись, разбор, дедуп.

    Возвращает событие, которое надо обработать, или `None` — «эту доставку уже
    обрабатывали, ответьте 2xx и не делайте ничего». Подпись не сошлась —
    `_BadSignatureError`.
    """
    if not verify_webhook(secret, body, headers.get("X-FastSub-Signature") or ""):
        raise _BadSignatureError

    import json as _json

    raw = _json.loads(body)
    delivery_id = headers.get("X-FastSub-Delivery-Id")
    if dedupe is not None and delivery_id and not dedupe(delivery_id):
        logger.info(
            "webhook %s: повтор доставки %s, пропускаем",
            raw.get("event"), delivery_id,
        )
        return None
    return _webhook_event(raw, delivery_id)


def _webhook_event(raw: dict[str, Any], delivery_id: str | None = None) -> WebhookEvent:
    name = raw.get("event") or ""
    cls = EVENT_CLASSES.get(name, WebhookEvent)
    return cls(
        event=raw.get("event") or "",
        link_id=raw.get("link_id") or "",
        status=raw.get("status") or "",
        user_id=raw.get("user_id"),
        task_id=raw.get("task_id"),
        publisher_payout_rub=str(raw.get("publisher_payout_rub") or "0"),
        payout_state=raw.get("payout_state") or "",
        payout_reversed=bool(raw.get("payout_reversed")),
        hold_until=raw.get("hold_until"),
        verified_at=raw.get("verified_at"),
        subscribed_at=raw.get("subscribed_at"),
        unsubscribed_at=raw.get("unsubscribed_at"),
        timestamp=raw.get("timestamp"),
        test=bool(raw.get("test")),
        delivery_id=delivery_id,
        raw=raw,
    )


def _invoke_sync(handler: Any, event: WebhookEvent) -> None:
    if isinstance(handler, WebhookEvents):
        handler.dispatch(event)
        return
    result = handler(event)
    if asyncio.iscoroutine(result):
        result.close()
        raise TypeError(
            "обработчик асинхронный, а приёмник синхронный — "
            "возьмите fastsub_webhook_router или fastsub_aiohttp_handler"
        )


def fastsub_webhook_router(
    secret: str,
    handler: Any,
    *,
    path: str = "/fastsub/webhook",
    dedupe: Any = True,
) -> Any:
    """Роутер FastAPI, принимающий наши события.

        app.include_router(fastsub_webhook_router(SECRET, on_event))

        async def on_event(event: WebhookEvent) -> None:
            if event.money_is_mine:
                ...

    Вместо функции можно передать `WebhookEvents` — тогда разбор по типам
    события возьмёт на себя он.

    Подпись проверяется до вызова `handler`, тело читается сырым. Ошибка внутри
    вашего обработчика превращается в 500 — мы повторим доставку позже.

    `dedupe` — что делать с повторами:

      * `True` (по умолчанию) — помнить `X-FastSub-Delivery-Id` в памяти
        процесса и второй раз обработчик не звать;
      * `False` — звать всегда, вы дедуплицируете сами (например, по `link_id`
        в своей БД, что надёжнее);
      * функция `(delivery_id: str) -> bool` — ваша проверка, `True` значит
        «видим впервые». Так подключается общий на все воркеры Redis:

            def seen(delivery_id: str) -> bool:
                return bool(redis.set(f"fs:{delivery_id}", 1, nx=True, ex=86400))
    """
    from fastapi import APIRouter, HTTPException, Request

    router = APIRouter()
    seen = _dedupe_fn(dedupe)

    async def _receive(request: Any) -> dict[str, bool]:
        try:
            event = _parse_delivery(secret, await request.body(), request.headers, seen)
        except _BadSignatureError:
            raise HTTPException(status_code=401, detail="bad signature") from None
        if event is not None:
            result = handler(event)
            if asyncio.iscoroutine(result):
                await result
        return {"ok": True}

    # Аннотацию проставляем классом, а не именем, и уже после определения.
    # В этом файле стоит `from __future__ import annotations`, поэтому все
    # аннотации — строки, а FastAPI разрешает их по глобальным именам модуля.
    # `Request` там нет и быть не может: fastapi импортируется внутри функции,
    # чтобы не стать обязательной зависимостью. Написанное декоратором
    # `request: Request` роняло `fastsub_webhook_router` прямо на создании
    # роутера — то есть приёмник для FastAPI не работал вообще ни у кого.
    _receive.__annotations__["request"] = Request
    router.post(path)(_receive)
    return router


def fastsub_aiohttp_handler(
    secret: str, handler: Any, *, dedupe: Any = True,
) -> Any:
    """То же для aiohttp: возвращает обработчик, который можно повесить на путь.

        app.router.add_post("/fastsub/webhook",
                            fastsub_aiohttp_handler(SECRET, on_event))

    `dedupe` работает так же, как в `fastsub_webhook_router`.
    """
    from aiohttp import web

    seen = _dedupe_fn(dedupe)

    async def _receive(request: Any) -> Any:
        try:
            event = _parse_delivery(secret, await request.read(), request.headers, seen)
        except _BadSignatureError:
            return web.json_response({"ok": False}, status=401)
        if event is not None:
            result = handler(event)
            if asyncio.iscoroutine(result):
                await result
        return web.json_response({"ok": True})

    return _receive


def fastsub_flask_blueprint(
    secret: str,
    handler: Any,
    *,
    path: str = "/fastsub/webhook",
    name: str = "fastsub",
    dedupe: Any = True,
) -> Any:
    """Blueprint для Flask — приёмник событий для синхронного бота.

        app.register_blueprint(fastsub_flask_blueprint(SECRET, on_event))

        def on_event(event: WebhookEvent) -> None:
            if event.money_was_taken:
                revoke_bonus(event.user_id)

    У асинхронных партнёров приёмник был, у синхронных — нет: те, кто взял
    `FastSubSync` и pyTelegramBotAPI, переписывали проверку подписи руками,
    ровно то, ради чего роутер для FastAPI и появился.

    Обработчик здесь обычная функция, не `async`. `WebhookEvents` тоже подойдёт,
    если его обработчики синхронные. `dedupe` — как в `fastsub_webhook_router`.
    """
    from flask import Blueprint, request

    bp = Blueprint(name, __name__)
    seen = _dedupe_fn(dedupe)

    def _receive() -> Any:
        try:
            # get_data() — сырое тело. get_json() пересобрал бы его, и подпись,
            # которая считается по байтам как есть, перестала бы сходиться.
            event = _parse_delivery(secret, request.get_data(), request.headers, seen)
        except _BadSignatureError:
            return {"ok": False}, 401
        if event is not None:
            _invoke_sync(handler, event)
        return {"ok": True}

    # Не декоратором: без установленного flask `bp.post` для mypy нетипизирован,
    # и строгий режим ругается на каждую функцию под ним.
    bp.add_url_rule(path, view_func=_receive, methods=["POST"])
    return bp


# ---------------------------------------------------------------------------
# Клиент рекламодателя
# ---------------------------------------------------------------------------
#
# Orders API существовал только в виде curl в документации: заказы, докупка,
# выгрузка подписчиков и оплата целевых действий делались руками.


class FastSubAdvertiser(_AsyncClient):
    """Заказы рекламодателя. Ключ — `fsa_live_…` из мини-аппа «Интеграция».

        async with FastSubAdvertiser(api_key="fsa_live_...") as adv:
            quote = await adv.quote(quantity=1000, price_rub="1.50")
            order = await adv.create_order(
                name="Мой канал", chat="@my_channel",
                quantity=1000, price_rub="1.50",
            )

    Транспорт общий с клиентом паблишера: те же повторы на 429 и 5xx, тот же
    `max_retries`. Отличается префикс пути — все методы живут под
    `/api/v1/advertiser`.
    """

    _prefix = "/api/v1/advertiser"
    _agent = " (advertiser)"

    #: Сколько секунд держать справочник таргетинга, не спрашивая заново.
    options_ttl = 600.0
    _options: dict[str, Any] | None = None
    _options_until: float = 0.0

    async def __aenter__(self) -> FastSubAdvertiser:
        return self

    @classmethod
    def from_env(
        cls,
        *,
        key_var: str = "FASTSUB_ADVERTISER_KEY",
        base_url_var: str = "FASTSUB_BASE_URL",
        **kwargs: Any,
    ) -> FastSubAdvertiser:
        """Клиент рекламодателя из переменных окружения — см. `FastSub.from_env`."""
        key = os.environ.get(key_var)
        if not key:
            raise FastSubError(
                0,
                f"переменная окружения {key_var} не задана — положите в неё "
                f"ключ вида fsa_live_… из мини-аппа «Интеграция»",
            )
        base_url = os.environ.get(base_url_var) or DEFAULT_BASE_URL
        return cls(api_key=key, base_url=base_url, **kwargs)

    async def targeting_options(self, *, fresh: bool = False) -> dict[str, Any]:
        """Справочник фильтров, коэффициентов и лимитов — тем же вызовом, что
        видит бот, чтобы не зашивать значения в свой код.

        Ответ кэшируется на `options_ttl` секунд (по умолчанию десять минут).
        Справочник меняется примерно раз в месяц, а дёргают его перед каждым
        расчётом цены — то есть в цикле по тысяче заказов это тысяча одинаковых
        ответов. `fresh=True` берёт свежий и обновляет кэш.
        """
        now = time.monotonic()
        if not fresh and self._options is not None and now < self._options_until:
            return self._options
        data = await self._send("GET", "/targeting/options")
        self._options = data
        self._options_until = now + self.options_ttl
        return data

    async def quote(self, **payload: Any) -> dict[str, Any]:
        """Сколько будет стоить. Ничего не списывает."""
        return await self._send("POST", "/orders/quote", body=payload)

    async def create_order(
        self, *, idempotency_key: str | None = None, **payload: Any
    ) -> dict[str, Any]:
        """Создать заказ и отправить на модерацию.

        `idempotency_key` — ваш уникальный ключ попытки. Повтор с тем же ключом
        вернёт уже созданный заказ, а не второй такой же.
        """
        headers = {"Idempotency-Key": idempotency_key} if idempotency_key else None
        return await self._send("POST", "/orders", body=payload, headers=headers)

    async def orders(self, **params: Any) -> dict[str, Any]:
        """Список заказов. Одна страница — всю сразу отдаёт `iter_orders`."""
        return await self._send("GET", "/orders", params=params)

    def iter_orders(
        self, *, page_size: int = 100, **params: Any,
    ) -> AsyncIterator[dict[str, Any]]:
        """Все заказы, страница за страницей.

            async for order in adv.iter_orders(status="active"):
                print(order["id"], order["progress"])
        """
        return _iter_pages(
            lambda path, *, params: self._send("GET", path, params=params),
            "/orders", params, page_size=page_size,
        )

    async def order(self, order_id: int) -> dict[str, Any]:
        """Карточка: прогресс, деньги, прогноз и удержание подписчиков."""
        return await self._send("GET", f"/orders/{order_id}")

    async def update_order(self, order_id: int, **payload: Any) -> dict[str, Any]:
        """Частичное изменение: пауза, цена, таргетинг, докупка, автопродление."""
        return await self._send("PATCH", f"/orders/{order_id}", body=payload)

    async def cancel_order(self, order_id: int) -> dict[str, Any]:
        """Отмена с возвратом неизрасходованного остатка."""
        return await self._send("DELETE", f"/orders/{order_id}")

    async def subscribers(self, order_id: int, **params: Any) -> dict[str, Any]:
        """Кто пришёл по заказу. Одна страница, без очереди задач."""
        return await self._send(
            "GET", f"/orders/{order_id}/subscribers", params=params,
        )

    def iter_subscribers(
        self, order_id: int, *, page_size: int = 100, **params: Any,
    ) -> AsyncIterator[dict[str, Any]]:
        """Все подписчики заказа, страница за страницей.

            async for sub in adv.iter_subscribers(order_id):
                await crm.upsert(sub["user_id"])

        Заказ на десять тысяч человек — это сотня страниц; лучше, чтобы их
        считал SDK, чем каждый партнёр заново.
        """
        return _iter_pages(
            lambda path, *, params: self._send("GET", path, params=params),
            f"/orders/{order_id}/subscribers", params, page_size=page_size,
        )

    async def postback(
        self,
        *,
        link_id: str,
        event: str,
        value_rub: str | float,
        meta: dict[str, Any] | None = None,
    ) -> dict[str, Any]:
        """Оплатить целевое действие: регистрацию, депозит, что угодно ваше.

        Идемпотентно по паре (выдача, событие): повтор вернёт `duplicate: true`
        и второй раз не спишет.
        """
        body: dict[str, Any] = {
            "link_id": link_id, "event": event, "value_rub": str(value_rub),
        }
        if meta:
            body["meta"] = meta
        return await self._send("POST", "/postback", body=body)

    async def confirm_start(self, *, start_param: str, user_id: int) -> dict[str, Any]:
        """Подтвердить запуск бота (задачи типа start_bot)."""
        return await self._send(
            "POST", "/confirm-start",
            body={"start_param": start_param, "user_id": user_id},
        )

    async def configure_webhook(
        self,
        url: str,
        *,
        events: list[str] | None = None,
        rotate_secret: bool = False,
    ) -> dict[str, Any]:
        """Куда слать события заказа. Секрет возвращается один раз — сохраните.

        До этого узнать, что стало с заказом, можно было только опросом
        `order(id)` в цикле: одобрили ночью, набрался к утру, встал в обед —
        и всё это вы видите, когда сами зайдёте посмотреть.

        События: order.approved, order.rejected, order.paused, order.resumed,
        order.completed, order.subscriber. Пусто — все.
        """
        body: dict[str, Any] = {"url": url, "rotate_secret": rotate_secret}
        if events:
            body["events"] = list(events)
        return await self._send("POST", "/webhook/configure", body=body)

    async def webhook_info(self) -> dict[str, Any]:
        """Что настроено сейчас. Секрет не отдаётся — он показан один раз."""
        return await self._send("GET", "/webhook")

    async def webhook_test(self) -> dict[str, Any]:
        """Фальшивый `order.approved` с полем `test: true`.

        Проверять приёмник на настоящем заказе — значит проверять его на живых
        деньгах: заказ одобряют один раз, второй попытки не будет.
        """
        return await self._send("POST", "/webhook/test", body={})

    async def start(self, message: Any) -> bool:
        """Подтвердить запуск прямо из хендлера /start. True — запуск засчитан.

            @dp.message(CommandStart())
            async def start(message: Message):
                await adv.start(message)          # ← вся интеграция
                await message.answer("Привет!")

        Разбирает `/start fastsub_…` сам: и параметр, и id юзера уже лежат в
        сообщении, а требовать доставать их руками — значит требовать прочитать,
        как устроен deep link, ради одного вызова.

        `False` — нет нашего параметра, запуск отклонён или API недоступен.
        Ваш бот работает как обычно; метод безопасен на каждом /start.

        Ничего не возвращает наружу и ничего не ломает: наш сбой не должен
        мешать юзеру запустить ваш бот. Повторный вызов с тем же параметром
        безопасен — подтверждение идемпотентно.
        """
        payload = start_payload(getattr(message, "text", None))
        user = getattr(message, "from_user", None)
        if payload is None or user is None:
            return False
        try:
            result = await self.confirm_start(start_param=payload, user_id=user.id)
        except FastSubError as exc:
            # Чужой параметр — не ошибка интеграции: так выглядит юзер, пришедший
            # мимо нас. А вот всё остальное стоит увидеть в логах.
            if exc.status not in (400, 404):
                logger.warning("confirm_start не удался: %s", exc)
            return False
        except Exception:
            logger.exception("непредвиденная ошибка confirm_start")
            return False
        return result.get("accepted") is True


#: Префикс, с которым мы выдаём deep link на бота рекламодателя.
START_PREFIX = "fastsub_"


def start_payload(text: str | None) -> str | None:
    """Наш параметр из текста `/start …`, либо None.

    Отдельной функцией, а не внутри `start()`: у telebot, Telegraf и голого
    вебхука сообщение устроено по-разному, а разбор один и тот же — и разбирать
    его руками в каждом боте значит четыре раза ошибиться в одном и том же.
    """
    parts = (text or "").strip().split(maxsplit=1)
    if len(parts) != 2 or not parts[0].startswith("/start"):
        return None
    payload = parts[1].strip()
    return payload if payload.startswith(START_PREFIX) else None



# ---------------------------------------------------------------------------
# Проверка интеграции: python -m fastsub doctor
# ---------------------------------------------------------------------------
#
# Половина обращений в поддержку — это четыре вопроса, на которые партнёр не
# может ответить сам: жив ли ключ, прошёл ли бот модерацию, включён ли веб-шаг,
# доходят ли вебхуки. Каждый из них отвечается одним запросом, но чтобы это
# узнать, надо сначала прочитать документацию — а пишут как раз те, кто до неё
# не дошёл.
#
# Поэтому команда, которая спрашивает всё это сама и отвечает по-русски. Ничего
# не чинит и ничего не меняет: диагностика, которая на ходу правит настройки, —
# это диагностика, после которой непонятно, что было сломано.

_OK = "  ok  "
_WARN = " ! "
_BAD = " ✗ "


@dataclass(slots=True)
class _Check:
    """Одна строка отчёта. `bad` — то, из-за чего интеграция не работает."""

    mark: str
    title: str
    detail: str = ""

    @property
    def bad(self) -> bool:
        return self.mark == _BAD

    def render(self) -> str:
        tail = f" — {self.detail}" if self.detail else ""
        return f"[{self.mark}] {self.title}{tail}"


def _check_account(fs: FastSubSync) -> list[_Check]:
    try:
        me = fs.me()
    except FastSubError as e:
        if e.status in (401, 403):
            return [_Check(_BAD, "Ключ", "не принят. Перевыпустите в мини-аппе «Интеграция»")]
        return [_Check(_BAD, "Ключ", f"проверить не удалось: {e.detail}")]

    checks = [
        _Check(_OK, "Ключ", f"аккаунт «{me.get('project_name') or '—'}»"),
        _Check(
            _OK, "Баланс",
            f"{me.get('balance_rub', '0')} ₽ доступно к выводу",
        ),
    ]
    hold = me.get("hold") or {}
    if hold:
        checks.append(_Check(_OK, "Отписки", str(hold.get("description") or "")))
    limits = (me.get("rate_limits") or {}).get("per_window") or {}
    if limits:
        checks.append(
            _Check(
                _OK, "Лимиты",
                ", ".join(f"{name} {value}/мин" for name, value in sorted(limits.items())),
            )
        )
    return checks


def _check_bot(fs: FastSubSync) -> list[_Check]:
    try:
        settings = fs.settings()
    except FastSubError as e:
        if e.status == 403 and "moderation" in e.detail:
            return [
                _Check(
                    _BAD, "Модерация",
                    "бот её ещё не прошёл: ключи работают, задания не выдаются",
                )
            ]
        if e.status == 403:
            return [_Check(_BAD, "Бот", e.detail)]
        if e.status == 401:
            return [_Check(_BAD, "Бот", "ключ не привязан к боту — перевыпустите его")]
        return [_Check(_BAD, "Бот", f"проверить не удалось: {e.detail}")]

    checks = [
        _Check(_OK, "Бот", f"«{settings.get('name')}», модерация пройдена"),
        _Check(
            _OK, "Выдача",
            f"{settings.get('sponsors_count')} спонсоров за запрос, "
            f"сброс списка раз в {int(settings.get('list_ttl_seconds', 0)) // 60} мин",
        ),
    ]
    if not settings.get("is_active"):
        checks.append(_Check(_BAD, "Бот выключен", "в карточке бота — заданий не будет"))

    if settings.get("show_quiz"):
        checks.append(_Check(_OK, "Веб-шаг", "анкета"))
    elif settings.get("use_smart_link"):
        checks.append(_Check(_OK, "Веб-шаг", "умный редирект"))
    else:
        checks.append(
            _Check(
                _WARN, "Веб-шаг", "выключен: заказы с гео- и демо-таргетингом недоступны",
            )
        )

    checks.append(
        _Check(
            _OK, "Блок ОП",
            "рисуете вы (get_links)" if settings.get("get_links") else "отправляем мы",
        )
    )
    floor = str(settings.get("min_reward_rub") or "0")
    if float(floor or 0) > 0:
        checks.append(
            _Check(_WARN, "Фильтр цены", f"не берём дешевле {floor} ₽ — поток меньше")
        )
    for name, label in (
        ("excluded_task_types", "типов заданий"),
        ("excluded_resource_types", "типов ресурсов"),
        ("excluded_themes", "тематик"),
    ):
        excluded = settings.get(name) or []
        if excluded:
            checks.append(
                _Check(_WARN, "Фильтр", f"исключено {len(excluded)} {label}: {', '.join(excluded)}")
            )
    return checks


def _check_webhook(fs: FastSubSync, *, send_test: bool) -> list[_Check]:
    try:
        info = fs.webhook_info()
    except FastSubError as e:
        return [_Check(_WARN, "Вебхук", f"проверить не удалось: {e.detail}")]

    url = info.get("url")
    if not url:
        return [_Check(_WARN, "Вебхук", "не настроен — статусы придётся опрашивать")]

    failures = int(info.get("consecutive_failures") or 0)
    mark = _BAD if failures >= 5 else (_WARN if failures else _OK)
    detail = str(url)
    if failures:
        detail += f", подряд неудач: {failures}"
    checks = [_Check(mark, "Вебхук", detail)]
    if send_test:
        try:
            fs.webhook_test()
            checks.append(_Check(_OK, "Тестовое событие", "отправлено, смотрите свой приёмник"))
        except FastSubError as e:
            checks.append(_Check(_BAD, "Тестовое событие", e.detail))
    return checks


def doctor(
    api_key: str | None = None,
    *,
    base_url: str | None = None,
    webhook_test: bool = False,
    out: Any = print,
) -> int:
    """Проверить интеграцию и напечатать отчёт. Возвращает код выхода.

        python -m fastsub doctor

    Ключ берётся из аргумента, иначе из `FASTSUB_KEY`. Ноль — всё, что мешает
    работать, в порядке; единица — есть хотя бы одна поломка. Предупреждения на
    код выхода не влияют: выключенный веб-шаг — это выбор, а не поломка.
    """
    key = api_key or os.environ.get("FASTSUB_KEY")
    if not key:
        out("Нужен ключ: python -m fastsub doctor <ключ> или FASTSUB_KEY в окружении.")
        return 2

    checks: list[_Check] = []
    with FastSubSync(api_key=key, base_url=base_url or DEFAULT_BASE_URL) as fs:
        out(f"FastSub {__version__} · {fs.base_url}\n")
        checks += _check_account(fs)
        # Ключ мёртв — остальные проверки ответят тем же 401 и утопят причину
        # в четырёх одинаковых строках.
        if not any(c.bad for c in checks):
            checks += _check_bot(fs)
            checks += _check_webhook(fs, send_test=webhook_test)

    for check in checks:
        out(check.render())

    broken = [c for c in checks if c.bad]
    out("")
    out("Всё в порядке." if not broken else f"Сломано: {len(broken)}.")
    return 1 if broken else 0


def _cli(argv: Sequence[str] | None = None) -> int:
    """`python -m fastsub …`. Одна команда, и та не требует ничего учить."""
    args = list(sys.argv[1:] if argv is None else argv)
    if not args or args[0] in ("-h", "--help", "help"):
        print(
            f"FastSub {__version__}\n\n"
            "  python -m fastsub doctor [ключ] [--webhook-test]\n"
            "      проверить интеграцию: ключ, бот, веб-шаг, вебхук, лимиты\n\n"
            "Ключ можно не передавать — возьмём из FASTSUB_KEY."
        )
        return 0
    if args[0] != "doctor":
        print(f"Неизвестная команда: {args[0]}. Есть только doctor.")
        return 2
    rest = args[1:]
    send_test = "--webhook-test" in rest
    positional = [a for a in rest if not a.startswith("-")]
    return doctor(
        positional[0] if positional else None,
        base_url=os.environ.get("FASTSUB_BASE_URL"),
        webhook_test=send_test,
    )
