|
1 | 1 | import asyncio |
2 | 2 | import datetime as dt |
3 | 3 | import uuid |
| 4 | +from typing import List, Tuple |
4 | 5 |
|
5 | 6 | import pytest |
6 | 7 | from taskiq import ScheduledTask |
7 | 8 |
|
8 | | -from taskiq_redis import RedisClusterScheduleSource, RedisScheduleSource |
| 9 | +from taskiq_redis import ( |
| 10 | + RedisClusterScheduleSource, |
| 11 | + RedisScheduleSource, |
| 12 | + RedisSentinelScheduleSource, |
| 13 | +) |
9 | 14 |
|
10 | 15 |
|
11 | 16 | @pytest.mark.anyio |
@@ -220,3 +225,140 @@ async def test_cluster_buffer(redis_cluster_url: str) -> None: |
220 | 225 | assert schedule1 in schedules |
221 | 226 | assert schedule2 in schedules |
222 | 227 | await source.shutdown() |
| 228 | + |
| 229 | + |
| 230 | +@pytest.mark.anyio |
| 231 | +async def test_sentinel_set_schedule( |
| 232 | + redis_sentinels: List[Tuple[str, int]], |
| 233 | + redis_sentinel_master_name: str, |
| 234 | +) -> None: |
| 235 | + prefix = uuid.uuid4().hex |
| 236 | + source = RedisSentinelScheduleSource( |
| 237 | + sentinels=redis_sentinels, |
| 238 | + master_name=redis_sentinel_master_name, |
| 239 | + prefix=prefix, |
| 240 | + ) |
| 241 | + schedule = ScheduledTask( |
| 242 | + task_name="test_task", |
| 243 | + labels={}, |
| 244 | + args=[], |
| 245 | + kwargs={}, |
| 246 | + cron="* * * * *", |
| 247 | + ) |
| 248 | + await source.add_schedule(schedule) |
| 249 | + schedules = await source.get_schedules() |
| 250 | + assert schedules == [schedule] |
| 251 | + await source.shutdown() |
| 252 | + |
| 253 | + |
| 254 | +@pytest.mark.anyio |
| 255 | +async def test_sentinel_delete_schedule( |
| 256 | + redis_sentinels: List[Tuple[str, int]], |
| 257 | + redis_sentinel_master_name: str, |
| 258 | +) -> None: |
| 259 | + prefix = uuid.uuid4().hex |
| 260 | + source = RedisSentinelScheduleSource( |
| 261 | + sentinels=redis_sentinels, |
| 262 | + master_name=redis_sentinel_master_name, |
| 263 | + prefix=prefix, |
| 264 | + ) |
| 265 | + schedule = ScheduledTask( |
| 266 | + task_name="test_task", |
| 267 | + labels={}, |
| 268 | + args=[], |
| 269 | + kwargs={}, |
| 270 | + cron="* * * * *", |
| 271 | + ) |
| 272 | + await source.add_schedule(schedule) |
| 273 | + schedules = await source.get_schedules() |
| 274 | + assert schedules == [schedule] |
| 275 | + await source.delete_schedule(schedule.schedule_id) |
| 276 | + schedules = await source.get_schedules() |
| 277 | + # Schedules are empty. |
| 278 | + assert not schedules |
| 279 | + await source.shutdown() |
| 280 | + |
| 281 | + |
| 282 | +@pytest.mark.anyio |
| 283 | +async def test_sentinel_post_run_cron( |
| 284 | + redis_sentinels: List[Tuple[str, int]], |
| 285 | + redis_sentinel_master_name: str, |
| 286 | +) -> None: |
| 287 | + prefix = uuid.uuid4().hex |
| 288 | + source = RedisSentinelScheduleSource( |
| 289 | + sentinels=redis_sentinels, |
| 290 | + master_name=redis_sentinel_master_name, |
| 291 | + prefix=prefix, |
| 292 | + ) |
| 293 | + schedule = ScheduledTask( |
| 294 | + task_name="test_task", |
| 295 | + labels={}, |
| 296 | + args=[], |
| 297 | + kwargs={}, |
| 298 | + cron="* * * * *", |
| 299 | + ) |
| 300 | + await source.add_schedule(schedule) |
| 301 | + assert await source.get_schedules() == [schedule] |
| 302 | + await source.post_send(schedule) |
| 303 | + assert await source.get_schedules() == [schedule] |
| 304 | + await source.shutdown() |
| 305 | + |
| 306 | + |
| 307 | +@pytest.mark.anyio |
| 308 | +async def test_sentinel_post_run_time( |
| 309 | + redis_sentinels: List[Tuple[str, int]], |
| 310 | + redis_sentinel_master_name: str, |
| 311 | +) -> None: |
| 312 | + prefix = uuid.uuid4().hex |
| 313 | + source = RedisSentinelScheduleSource( |
| 314 | + sentinels=redis_sentinels, |
| 315 | + master_name=redis_sentinel_master_name, |
| 316 | + prefix=prefix, |
| 317 | + ) |
| 318 | + schedule = ScheduledTask( |
| 319 | + task_name="test_task", |
| 320 | + labels={}, |
| 321 | + args=[], |
| 322 | + kwargs={}, |
| 323 | + time=dt.datetime(2000, 1, 1), |
| 324 | + ) |
| 325 | + await source.add_schedule(schedule) |
| 326 | + assert await source.get_schedules() == [schedule] |
| 327 | + await source.post_send(schedule) |
| 328 | + assert await source.get_schedules() == [] |
| 329 | + await source.shutdown() |
| 330 | + |
| 331 | + |
| 332 | +@pytest.mark.anyio |
| 333 | +async def test_sentinel_buffer( |
| 334 | + redis_sentinels: List[Tuple[str, int]], |
| 335 | + redis_sentinel_master_name: str, |
| 336 | +) -> None: |
| 337 | + prefix = uuid.uuid4().hex |
| 338 | + source = RedisSentinelScheduleSource( |
| 339 | + sentinels=redis_sentinels, |
| 340 | + master_name=redis_sentinel_master_name, |
| 341 | + prefix=prefix, |
| 342 | + buffer_size=1, |
| 343 | + ) |
| 344 | + schedule1 = ScheduledTask( |
| 345 | + task_name="test_task1", |
| 346 | + labels={}, |
| 347 | + args=[], |
| 348 | + kwargs={}, |
| 349 | + cron="* * * * *", |
| 350 | + ) |
| 351 | + schedule2 = ScheduledTask( |
| 352 | + task_name="test_task2", |
| 353 | + labels={}, |
| 354 | + args=[], |
| 355 | + kwargs={}, |
| 356 | + cron="* * * * *", |
| 357 | + ) |
| 358 | + await source.add_schedule(schedule1) |
| 359 | + await source.add_schedule(schedule2) |
| 360 | + schedules = await source.get_schedules() |
| 361 | + assert len(schedules) == 2 |
| 362 | + assert schedule1 in schedules |
| 363 | + assert schedule2 in schedules |
| 364 | + await source.shutdown() |
0 commit comments