Skip to content

Commit 79c3440

Browse files
committed
fix(webhook): устранить гонки lifecycle
Запрос во время shutdown: Начальная проверка могла пройти до route.match(), security, resolve bot или JSON, после чего TaskTracker уже закрывался, а update запускался и получал ответ 200. Состояние теперь проверяется повторно прямо перед spawn без последующих await. Общий guard выбран, чтобы одинаково защитить все webhook engines. Конкурентные startup и shutdown: Startup мог снять флаг остановки до завершения параллельного shutdown. Lifecycle callbacks сериализованы одним asyncio.Lock. Это сохраняет повторный startup, но исключает наложение переходов состояния. Подписка SingleBotEngine до startup: Route с BotIdParam обращался к bot.info до загрузки данных и падал StateError. Subscribe теперь вызывает get_my_info только для еще не запущенного бота. Проверка стоит до webhook kwargs и URL, потому что оба могут зависеть от info. Частичный startup нескольких ботов: Ошибка второго бота оставляла первый запущенным, а workflow_data содержал stale bot. Аварийный путь теперь вызывает shutdown-сигналы, закрывает новые экземпляры, восстанавливает реестр и workflow_data. Rollback повторяет lifecycle-контракт и не оставляет частично активное состояние. Порядок multi-bot shutdown: BeforeShutdown вызывался до завершения фоновых update handlers. TaskTracker теперь сначала дожидается задач, затем идут shutdown-сигнал и close. Порядок выбран по уже используемому SingleBotEngine и исключает гонку cleanup. Конкурентный реестр TokenEngine: Add, remove, resolve и shutdown могли одновременно менять bots и token_ids, терять Bot, создавать tracker для удаленного бота или менять dict при обходе. Один asyncio.Lock сериализует редкие операции с реестром, а identity-check запрещает spawn для уже удаленного или замененного экземпляра. Один lock проще per-bot locks и достаточен без измеренной проблемы throughput. Замена и удаление Bot: Старый mapping удалялся до успешного unsubscribe, а новый кандидат мог утечь при ошибке cleanup. Mapping теперь меняется только после успешного unsubscribe, tracker завершается до close, а кандидаты закрываются на обоих путях add и request resolve. Так ошибка сохраняет рабочий старый Bot и не оставляет локальные ресурсы. Ошибки Bot API при resolve: Любой MaxBotApiError превращался в 404, включая временную недоступность MAX. Ошибки токена по-прежнему скрываются за 404, а серверные ошибки оборачиваются в BotStartError и возвращают контролируемый 502. Это отличает неизвестного бота от сбоя внешнего API без раскрытия токена. Проверки: Добавлены регрессионные тесты для гонок, rollback, BotIdParam, malformed JSON, ошибок API и отказов cleanup. Полный прогон: 1705 тестов, ruff, mypy, codespell, slotscheck и bandit.
1 parent 5437eb4 commit 79c3440

9 files changed

Lines changed: 622 additions & 79 deletions

File tree

src/maxo/transport/webhook/engines/base.py

Lines changed: 15 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,4 @@
1+
import asyncio
12
from abc import ABC, abstractmethod
23
from typing import Any, Generic, TypeVar
34

@@ -44,6 +45,7 @@ def __init__(
4445

4546
self.shutdown_timeout = shutdown_timeout
4647
self._is_shutting_down = False
48+
self._lifecycle_lock = asyncio.Lock()
4749

4850
def register(self, app: AppT) -> None:
4951
webhook.info(
@@ -71,8 +73,7 @@ async def handle_request(
7173
request: WebRequest[RawRequestT],
7274
) -> FrameworkResponseT:
7375
try:
74-
if self._is_shutting_down:
75-
raise RequestHandlingStoppedError
76+
self._ensure_accepting_requests()
7677

7778
route_params = await self.route.match(request)
7879

@@ -98,6 +99,8 @@ async def handle_request(
9899

99100
webhook.debug("New update: %s", update.update)
100101

102+
self._ensure_accepting_requests(bot)
103+
101104
self._get_task_tracker(bot).spawn( # type: ignore[unused-awaitable]
102105
self._feed_update(bot=bot, update=update),
103106
)
@@ -121,12 +124,14 @@ async def handle_request(
121124
)
122125

123126
async def on_startup(self, app: AppT, *args: Any, **kwargs: Any) -> None:
124-
await self._on_startup(app, *args, **kwargs)
125-
self._is_shutting_down = False
127+
async with self._lifecycle_lock:
128+
await self._on_startup(app, *args, **kwargs)
129+
self._is_shutting_down = False
126130

127131
async def on_shutdown(self, app: AppT, *args: Any, **kwargs: Any) -> None:
128-
self._is_shutting_down = True
129-
await self._on_shutdown(app, *args, **kwargs)
132+
async with self._lifecycle_lock:
133+
self._is_shutting_down = True
134+
await self._on_shutdown(app, *args, **kwargs)
130135

131136
@abstractmethod
132137
async def _on_startup(self, app: AppT, *args: Any, **kwargs: Any) -> None:
@@ -171,3 +176,7 @@ async def _feed_update(self, bot: Bot, update: MaxoUpdate[Any]) -> None:
171176
result = await self.dispatcher.feed_update(bot=bot, update=update)
172177
if isinstance(result, MaxoMethod):
173178
await bot.silent_call_method(method=result)
179+
180+
def _ensure_accepting_requests(self, bot: Bot | None = None) -> None:
181+
if self._is_shutting_down or (bot is not None and bot.closed):
182+
raise RequestHandlingStoppedError

src/maxo/transport/webhook/engines/multi.py

Lines changed: 51 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818
FrameworkResponseT,
1919
RawRequestT,
2020
)
21+
from maxo.transport.webhook.engines.errors import RequestHandlingStoppedError
2122
from maxo.transport.webhook.route import Route
2223
from maxo.transport.webhook.security import Security
2324
from maxo.transport.webhook.tasks import TaskTracker
@@ -60,7 +61,10 @@ async def _on_startup(
6061
bots: Iterable[Bot] | None = None,
6162
**kwargs: Any,
6263
) -> None:
64+
initial_bots = dict(self._bots)
65+
initial_workflow_data = dict(self.dispatcher.workflow_data)
6366
all_bots = set(self.bots.values())
67+
processed_bots: list[Bot] = []
6468

6569
if bots is not None:
6670
all_bots |= set(bots)
@@ -69,16 +73,47 @@ async def _on_startup(
6973
"Starting multi-bot webhook engine with %s bot(s)",
7074
len(all_bots),
7175
)
72-
for bot in all_bots:
73-
lifecycle_data = self._build_lifecycle_data(app=app, bot=bot, **kwargs)
74-
self.dispatcher.workflow_data.update(lifecycle_data)
76+
try:
77+
for bot in all_bots:
78+
processed_bots.append(bot)
79+
lifecycle_data = self._build_lifecycle_data(
80+
app=app,
81+
bot=bot,
82+
**kwargs,
83+
)
84+
self.dispatcher.workflow_data.update(lifecycle_data)
7585

76-
await self.dispatcher.feed_signal(BeforeStartup(), bot)
77-
await bot.get_my_info()
78-
self._bots[bot.info.id] = bot
79-
await self.dispatcher.feed_signal(AfterStartup(), bot)
86+
await self.dispatcher.feed_signal(BeforeStartup(), bot)
87+
await bot.get_my_info()
88+
self._bots[bot.info.id] = bot
89+
await self.dispatcher.feed_signal(AfterStartup(), bot)
90+
except BaseException:
91+
try:
92+
if processed_bots:
93+
lifecycle_data = self._build_lifecycle_data(
94+
app=app,
95+
bots=set(processed_bots),
96+
**kwargs,
97+
)
98+
self.dispatcher.workflow_data.update(lifecycle_data)
99+
await self.dispatcher.feed_signal(BeforeShutdown())
100+
101+
for bot in all_bots - set(initial_bots.values()):
102+
await bot.close()
103+
104+
if processed_bots:
105+
await self.dispatcher.feed_signal(AfterShutdown())
106+
finally:
107+
self._bots.clear()
108+
self._bots.update(initial_bots)
109+
self.dispatcher.workflow_data.clear()
110+
self.dispatcher.workflow_data.update(initial_workflow_data)
111+
raise
80112

81113
def _get_task_tracker(self, bot: Bot) -> TaskTracker:
114+
if self._bots.get(bot.info.id) is not bot:
115+
raise RequestHandlingStoppedError
116+
82117
tracker = self._task_trackers.get(bot.info.id)
83118

84119
if tracker is None:
@@ -88,24 +123,28 @@ def _get_task_tracker(self, bot: Bot) -> TaskTracker:
88123
return tracker
89124

90125
async def _on_shutdown(self, app: AppT, *args: Any, **kwargs: Any) -> None:
126+
bots = tuple(self._bots.values())
91127
webhook.info(
92128
"Stopping %s with %s bot(s)",
93129
self.__class__.__name__,
94-
len(self._bots),
130+
len(bots),
95131
)
96132

133+
for bot in bots:
134+
if (tracker := self._task_trackers.pop(bot.info.id, None)) is not None:
135+
await tracker.close(timeout=self.shutdown_timeout)
136+
97137
lifecycle_data = self._build_lifecycle_data(
98138
app=app,
99-
bots=set(self._bots.values()),
139+
bots=set(bots),
100140
**kwargs,
101141
)
102142
self.dispatcher.workflow_data.update(lifecycle_data)
103143
await self.dispatcher.feed_signal(BeforeShutdown())
104144

105-
for bot in self._bots.values():
106-
if (tracker := self._task_trackers.pop(bot.info.id, None)) is not None:
107-
await tracker.close(timeout=self.shutdown_timeout)
145+
for bot in bots:
108146
await bot.close()
109147

110148
self._bots.clear()
149+
self._task_trackers.clear()
111150
await self.dispatcher.feed_signal(AfterShutdown())

src/maxo/transport/webhook/engines/single.py

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -58,6 +58,9 @@ async def subscribe(
5858
self,
5959
webhook_config: WebhookConfig | None = None,
6060
) -> SimpleQueryResult:
61+
if not self.bot.started:
62+
await self.bot.get_my_info()
63+
6164
kwargs = await self._build_webhook_kwargs(
6265
bot=self.bot,
6366
base_config=webhook_config or WebhookConfig(),
Lines changed: 102 additions & 59 deletions
Original file line numberDiff line numberDiff line change
@@ -1,11 +1,19 @@
1+
import asyncio
12
from typing import Any, Generic
23

34
from maxo import Bot, Dispatcher
4-
from maxo.errors import MaxBotApiError
5+
from maxo.errors import (
6+
MaxBotApiError,
7+
MaxBotBadRequestError,
8+
MaxBotForbiddenError,
9+
MaxBotNotFoundError,
10+
MaxBotUnauthorizedError,
11+
)
512
from maxo.loggers import webhook
613
from maxo.transport.webhook.configs.bot import BotConfig
714
from maxo.transport.webhook.configs.webhook import WebhookConfig
815
from maxo.transport.webhook.engines.base import AppT, FrameworkResponseT, RawRequestT
16+
from maxo.transport.webhook.engines.errors import BotStartError
917
from maxo.transport.webhook.engines.multi import BaseMultiBotEngine
1018
from maxo.transport.webhook.route import Route
1119
from maxo.transport.webhook.route.params import RouteParams
@@ -43,84 +51,119 @@ def __init__(
4351
)
4452
self.bot_config = bot_config
4553
self._token_ids: dict[str, int] = {}
54+
# ponytail: один lock на реестр, per-bot locks нужны только при замерах задержек
55+
self._bots_lock = asyncio.Lock()
4656

4757
async def add_bot(
4858
self,
4959
token: str,
5060
webhook_config: WebhookConfig | None = None,
5161
) -> Bot:
52-
bot = self._build_bot(token)
53-
await bot.get_my_info()
62+
async with self._bots_lock:
63+
bot = self._build_bot(token)
64+
await bot.get_my_info()
5465

55-
kwargs = await self._build_webhook_kwargs(
56-
bot=bot,
57-
base_config=webhook_config or self.webhook_config,
58-
)
59-
await bot.subscribe(url=await self.route.build_url(bot=bot), **kwargs)
66+
kwargs = await self._build_webhook_kwargs(
67+
bot=bot,
68+
base_config=webhook_config or self.webhook_config,
69+
)
70+
await bot.subscribe(url=await self.route.build_url(bot=bot), **kwargs)
71+
72+
if (existing := self._bots.get(bot.info.id)) is not None:
73+
# URL contains the token: same token = same URL, already subscribed
74+
try:
75+
await self._remove_bot(
76+
bot.info.id,
77+
unsubscribe=existing.token != token,
78+
)
79+
except Exception:
80+
await bot.close()
81+
raise
82+
83+
self._bots[bot.info.id] = bot
84+
self._token_ids[token] = bot.info.id
85+
86+
webhook.info("Added bot %s to token engine and set webhook", bot.info.id)
87+
return bot
6088

61-
if (existing := self._bots.get(bot.info.id)) is not None:
62-
# URL contains the token: same token = same URL, already subscribed
63-
await self.remove_bot(bot.info.id, unsubscribe=existing.token != token)
89+
async def remove_bot(self, bot_id: int, unsubscribe: bool = True) -> bool:
90+
async with self._bots_lock:
91+
return await self._remove_bot(bot_id, unsubscribe=unsubscribe)
6492

65-
self._bots[bot.info.id] = bot
66-
self._token_ids[token] = bot.info.id
93+
async def _resolve_bot(self, route_params: RouteParams) -> Bot | None:
94+
async with self._bots_lock:
95+
token = route_params.get("bot_token")
96+
if not isinstance(token, str) or not token:
97+
return None
98+
99+
if (bot_id := self._token_ids.get(token)) is not None:
100+
return self._bots.get(bot_id)
101+
102+
bot = self._build_bot(token)
103+
try:
104+
await bot.get_my_info()
105+
except (
106+
MaxBotBadRequestError,
107+
MaxBotForbiddenError,
108+
MaxBotNotFoundError,
109+
MaxBotUnauthorizedError,
110+
):
111+
await bot.close()
112+
return None
113+
except MaxBotApiError as exc:
114+
await bot.close()
115+
raise BotStartError(original_error=exc) from exc
116+
117+
existing = self._bots.get(bot.info.id)
118+
if existing is not None and existing.token == token:
119+
self._token_ids[token] = existing.info.id
120+
await bot.close()
121+
return existing
122+
123+
if existing is not None:
124+
try:
125+
await self._remove_bot(existing.info.id)
126+
except Exception as exc:
127+
await bot.close()
128+
raise BotStartError(original_error=exc) from exc
129+
130+
self._bots[bot.info.id] = bot
131+
self._token_ids[token] = bot.info.id
132+
return bot
67133

68-
webhook.info("Added bot %s to token engine and set webhook", bot.info.id)
69-
return bot
134+
async def _on_shutdown(self, app: AppT, *args: Any, **kwargs: Any) -> None:
135+
async with self._bots_lock:
136+
await super()._on_shutdown(app, *args, **kwargs)
137+
self._token_ids.clear()
70138

71-
async def remove_bot(self, bot_id: int, unsubscribe: bool = True) -> bool:
72-
bot = self._bots.pop(bot_id, None)
139+
def _build_bot(self, token: str) -> Bot:
140+
return Bot(
141+
token=token,
142+
client=self.bot_config.client,
143+
defaults=self.bot_config.defaults,
144+
upload_config=self.bot_config.upload_config,
145+
warming_up=self.bot_config.warming_up,
146+
)
147+
148+
async def _remove_bot(self, bot_id: int, *, unsubscribe: bool = True) -> bool:
149+
bot = self._bots.get(bot_id)
73150
if bot is None:
74151
return False
75152

76-
self._token_ids.pop(bot.token, None)
153+
if unsubscribe:
154+
await bot.unsubscribe(url=await self.route.build_url(bot=bot))
77155

156+
self._bots.pop(bot_id, None)
157+
self._token_ids.pop(bot.token, None)
78158
try:
79-
if unsubscribe:
80-
await bot.unsubscribe(url=await self.route.build_url(bot=bot))
81-
finally:
82159
if (tracker := self._task_trackers.pop(bot_id, None)) is not None:
83160
await tracker.close(timeout=self.shutdown_timeout)
161+
finally:
84162
await bot.close()
85163

86164
webhook.info("Removed bot %s from token engine", bot_id)
87165
return True
88166

89-
async def _resolve_bot(self, route_params: RouteParams) -> Bot | None:
90-
token = route_params.get("bot_token")
91-
if not isinstance(token, str) or not token:
92-
return None
93-
94-
if (bot_id := self._token_ids.get(token)) is not None:
95-
return self._bots.get(bot_id)
96-
97-
bot = self._build_bot(token)
98-
try:
99-
await bot.get_my_info()
100-
except MaxBotApiError:
101-
return None
102-
103-
existing = self._bots.get(bot.info.id)
104-
if existing is not None and existing.token == token:
105-
self._token_ids[token] = existing.info.id
106-
return existing
107-
108-
if existing is not None:
109-
self._token_ids.pop(existing.token, None)
110-
111-
self._bots[bot.info.id] = bot
112-
self._token_ids[token] = bot.info.id
113-
return bot
114-
115-
async def _on_shutdown(self, app: AppT, *args: Any, **kwargs: Any) -> None:
116-
await super()._on_shutdown(app, *args, **kwargs)
117-
self._token_ids.clear()
118-
119-
def _build_bot(self, token: str) -> Bot:
120-
return Bot(
121-
token=token,
122-
client=self.bot_config.client,
123-
defaults=self.bot_config.defaults,
124-
upload_config=self.bot_config.upload_config,
125-
warming_up=self.bot_config.warming_up,
126-
)
167+
async def _on_startup(self, app: AppT, *args: Any, **kwargs: Any) -> None:
168+
async with self._bots_lock:
169+
await super()._on_startup(app, *args, **kwargs)

tests/maxo_webhook/fixtures/shutdown.py

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -85,3 +85,22 @@ async def make_request(self, request: HTTPRequest) -> HTTPResponse[Any]:
8585

8686
async def close(self) -> None:
8787
self.closed = True
88+
89+
90+
class SubscribingClient(TrackableClient):
91+
def __init__(self, bot_id: int = 42) -> None:
92+
super().__init__(bot_id=bot_id)
93+
self.subscribed: list[str] = []
94+
95+
async def make_request(self, request: HTTPRequest) -> HTTPResponse[Any]:
96+
if request.url != "subscriptions":
97+
return await super().make_request(request)
98+
99+
self.subscribed.append(request.body["url"])
100+
return HTTPResponse(
101+
status_code=200,
102+
headers={},
103+
data={"success": True},
104+
cookies={},
105+
raw_response=None,
106+
)

0 commit comments

Comments
 (0)