Skip to content

Commit a2e1f9f

Browse files
committed
chore(feat): add scheduler
1 parent f598e88 commit a2e1f9f

11 files changed

Lines changed: 214 additions & 4 deletions

File tree

pyproject.toml

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ requires-python = ">=3.12"
77
dependencies = [
88
"aiogram>=3.26.0",
99
"alembic>=1.18.4",
10+
"apscheduler>=3.11.2",
1011
"asyncpg>=0.31.0",
1112
"dependency-injector2>=4.41.1",
1213
"mashumaro>=3.20",
@@ -31,7 +32,7 @@ known-first-party = ["src"]
3132

3233

3334
[[tool.mypy.overrides]]
34-
module = ["vkbottle.*", "vk_api.*", "vk_api"]
35+
module = ["vkbottle.*", "vk_api.*", "vk_api","apscheduler.*"]
3536
ignore_missing_imports = true
3637
[tool.alembic]
3738

src/domain/repositories/user_repository.py

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
from abc import ABC, abstractmethod
2-
from typing import List
2+
from typing import AsyncIterator, List
33

44
from domain.dtos import SaveUserDTO
55
from domain.entities import User
@@ -20,6 +20,9 @@ async def get_by_external_id(
2020
@abstractmethod
2121
async def get_users_for_reminder(self, hour: int) -> List[User]: ...
2222

23+
@abstractmethod
24+
def iter_users_for_reminder(self, hour: int) -> AsyncIterator[User]: ...
25+
2326
# @abstractmethod
2427
# def find_by_id(self, user_id: int) -> Optional[User]:
2528
# """Find user by ID"""

src/infrastructure/database/repositories/sqlachemy/user_repository.py

Lines changed: 32 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
import logging
2-
from typing import List
2+
from typing import AsyncIterator, List
33

44
from sqlalchemy import select
55
from sqlalchemy.exc import IntegrityError
@@ -16,6 +16,8 @@
1616

1717

1818
class SQLAchemyUserRepository(UserRepository):
19+
BATCH_SIZE = 100
20+
1921
def __init__(self, session_manager: DatabaseSessionManager):
2022
self.async_session_maker = session_manager
2123

@@ -60,6 +62,35 @@ async def get_by_external_id(
6062

6163
return self._model_to_entity(user_model)
6264

65+
async def iter_users_for_reminder(self, hour: int) -> AsyncIterator[User]:
66+
offset = 0
67+
while True:
68+
async with self.async_session_maker.get_session() as session:
69+
stmt = (
70+
select(UserModel)
71+
.where(
72+
UserModel.reminder_enabled.is_(True),
73+
UserModel.reminder_hour == hour,
74+
UserModel.deleted_at.is_(None),
75+
)
76+
.order_by(UserModel.id)
77+
.limit(self.BATCH_SIZE)
78+
.offset(offset)
79+
)
80+
result = await session.execute(stmt)
81+
models = result.scalars().all()
82+
83+
if not models:
84+
break
85+
86+
for model in models:
87+
yield self._model_to_entity(model)
88+
89+
offset += self.BATCH_SIZE
90+
91+
if len(models) < self.BATCH_SIZE:
92+
break
93+
6394
def _model_to_entity(self, model: UserModel) -> User:
6495
return User(
6596
id=model.id,

src/infrastructure/ioc/container/infrastructure.py

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@
1111
SQLAchemyDiaryRepository,
1212
SQLAchemyUserRepository,
1313
)
14+
from infrastructure.scheduler.scheduler import AppScheduler
1415

1516

1617
class InfrastructureContainer(containers.DeclarativeContainer):
@@ -42,3 +43,7 @@ class InfrastructureContainer(containers.DeclarativeContainer):
4243
chart_generator: Factory[ChartGeneratorInterface] = providers.Factory(
4344
MoodChartGenerator,
4445
)
46+
47+
scheduler: Singleton[AppScheduler] = providers.Singleton(
48+
AppScheduler,
49+
)
Lines changed: 77 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,77 @@
1+
import logging
2+
from typing import Callable, Coroutine, Any
3+
4+
from apscheduler.schedulers.asyncio import AsyncIOScheduler
5+
from apscheduler.triggers.cron import CronTrigger
6+
from apscheduler.triggers.interval import IntervalTrigger
7+
8+
logger = logging.getLogger(__name__)
9+
10+
11+
class AppScheduler:
12+
def __init__(self, timezone: str = "Europe/Moscow"):
13+
self.scheduler = AsyncIOScheduler(
14+
timezone=timezone,
15+
job_defaults={
16+
"coalesce": True,
17+
"max_instances": 2,
18+
"misfire_grace_time": 300,
19+
},
20+
)
21+
22+
def add_cron_job(
23+
self,
24+
func: Callable[..., Coroutine[Any, Any, None]],
25+
*,
26+
id: str,
27+
name: str | None = None,
28+
hour: int | str = "*",
29+
minute: int | str = "*",
30+
day_of_week: str = "*",
31+
args: tuple | None = None,
32+
kwargs: dict | None = None,
33+
) -> None:
34+
self.scheduler.add_job(
35+
func,
36+
trigger=CronTrigger(
37+
hour=hour,
38+
minute=minute,
39+
day_of_week=day_of_week,
40+
),
41+
id=id,
42+
name=name or id,
43+
args=args or (),
44+
kwargs=kwargs or {},
45+
)
46+
logger.info("Registered cron job: %s (%s:%s)", id, hour, minute)
47+
48+
def add_interval_job(
49+
self,
50+
func: Callable[..., Coroutine[Any, Any, None]],
51+
*,
52+
id: str,
53+
name: str | None = None,
54+
minutes: int = 1,
55+
args: tuple | None = None,
56+
kwargs: dict | None = None,
57+
) -> None:
58+
self.scheduler.add_job(
59+
func,
60+
trigger=IntervalTrigger(minutes=minutes),
61+
id=id,
62+
name=name or id,
63+
args=args or (),
64+
kwargs=kwargs or {},
65+
)
66+
logger.info("Registered interval job: %s (every %d min)", id, minutes)
67+
68+
async def start(self) -> None:
69+
self.scheduler.start()
70+
logger.info(
71+
"Scheduler started with %d jobs", len(self.scheduler.get_jobs())
72+
)
73+
74+
async def shutdown(self) -> None:
75+
if self.scheduler.running:
76+
self.scheduler.shutdown(wait=True)
77+
logger.info("Scheduler stopped gracefully")

src/main.py

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,9 +32,11 @@
3232

3333
async def on_startup() -> None:
3434
await container.infrastructure.container.cache().get_connection()
35+
3536
session_manager = container.infrastructure.container.session_manager()
3637
async with session_manager.get_session() as session:
3738
await session.execute(text("SELECT 1"))
39+
3840
logger.info("All connections warmed up")
3941

4042

@@ -43,6 +45,7 @@ async def on_shutdown() -> None:
4345
await executor_pool.shutdown_all()
4446
await container.infrastructure.redis_cache().close()
4547
await container.infrastructure.session_manager().close()
48+
await container.infrastructure.scheduler().shutdown()
4649
logger.info("Cleanup completed")
4750

4851

@@ -83,6 +86,8 @@ async def async_main() -> None:
8386
signal_handler.install_handlers()
8487

8588
runner = BotRunner(bots)
89+
90+
await container.infrastructure.scheduler().start()
8691

8792
await asyncio.gather(
8893
runner.start_all(),

src/presentation/common/messages.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -257,6 +257,7 @@ def get_profile_text_with_stats(
257257

258258
REMINDER_ENABLE_TEXT = "⏰ Включить напоминания"
259259
REMINDER_EDIT_TEXT = "Изменить время напоминания {current}"
260+
REMINDER_TEXT = "🔔 Напоминание: отметить настроение за сегодня!"
260261

261262
INVALID_PERIOD = "❌ Неверный период"
262263

src/presentation/vk/bot.py

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@
44

55
from presentation.common.base_bot import BaseBot
66
from presentation.vk.handlers.router import VkRouter, create_vk_router
7+
from presentation.vk.notifications.mood_record_reminder import MoodRecordReminder
78
from presentation.vk.polling import VkLongPolling
89
from presentation.vk.sdk.api import VkSdk
910
from presentation.vk.sdk.types import VkMessage
@@ -58,7 +59,14 @@ async def start(self) -> None:
5859
vk = vk_api.VkApi(token=self._token)
5960
vk_sdk = VkSdk(vk)
6061
self._router = create_vk_router(vk_sdk, self._container, self._group_id)
61-
62+
63+
mood_record_reminder = MoodRecordReminder(
64+
vk_api=vk_sdk,
65+
user_repository=self._container.infrastructure.container.user_repository(),
66+
scheduler=self._container.infrastructure.scheduler(),
67+
)
68+
await mood_record_reminder.register()
69+
6270
self._polling = VkLongPolling(
6371
token=self._token,
6472
group_id=self._group_id,

src/presentation/vk/notifications/__init__.py

Whitespace-only changes.
Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,44 @@
1+
import logging
2+
import time
3+
4+
from domain.entities import User
5+
from domain.repositories import UserRepository
6+
from infrastructure.scheduler.scheduler import AppScheduler
7+
from presentation.common import Messages
8+
from presentation.vk.sdk.api import VkSdk
9+
10+
logger = logging.getLogger(__name__)
11+
12+
13+
class MoodRecordReminder:
14+
def __init__(
15+
self, vk_api: VkSdk, user_repository: UserRepository, scheduler: AppScheduler
16+
) -> None:
17+
self.vk_api = vk_api
18+
self._user_repository = user_repository
19+
self._scheduler = scheduler
20+
21+
async def register(self) -> None:
22+
self._scheduler.add_cron_job(
23+
self._notify,
24+
id="hourly_reminder",
25+
name="hourly_reminder",
26+
hour="*",
27+
minute="0",
28+
)
29+
30+
async def _notify(self) -> None:
31+
hour = time.localtime().tm_hour
32+
processed = 0
33+
async for user in self._user_repository.iter_users_for_reminder(hour):
34+
try:
35+
await self._send_reminder(user)
36+
processed += 1
37+
except Exception:
38+
logger.exception("Failed to send reminder to user %s", user.id)
39+
logger.info(
40+
"Reminder job completed: %d users processed for hour=%d", processed, hour
41+
)
42+
43+
async def _send_reminder(self, user: User) -> None:
44+
await self.vk_api.send_message(int(user.external_id), Messages.REMINDER_TEXT)

0 commit comments

Comments
 (0)