|
| 1 | +import datetime as dt |
1 | 2 | import uuid
|
2 | 3 |
|
3 | 4 | import pytest
|
4 | 5 | from taskiq import ScheduledTask
|
5 | 6 |
|
6 |
| -from taskiq_redis import RedisScheduleSource |
| 7 | +from taskiq_redis import RedisClusterScheduleSource, RedisScheduleSource |
7 | 8 |
|
8 | 9 |
|
9 | 10 | @pytest.mark.anyio
|
@@ -56,6 +57,153 @@ async def test_post_run_cron(redis_url: str) -> None:
|
56 | 57 | cron="* * * * *",
|
57 | 58 | )
|
58 | 59 | await source.add_schedule(schedule)
|
| 60 | + assert await source.get_schedules() == [schedule] |
| 61 | + await source.post_send(schedule) |
| 62 | + assert await source.get_schedules() == [schedule] |
| 63 | + await source.shutdown() |
| 64 | + |
| 65 | + |
| 66 | +@pytest.mark.anyio |
| 67 | +async def test_post_run_time(redis_url: str) -> None: |
| 68 | + prefix = uuid.uuid4().hex |
| 69 | + source = RedisScheduleSource(redis_url, prefix=prefix) |
| 70 | + schedule = ScheduledTask( |
| 71 | + task_name="test_task", |
| 72 | + labels={}, |
| 73 | + args=[], |
| 74 | + kwargs={}, |
| 75 | + time=dt.datetime(2000, 1, 1), |
| 76 | + ) |
| 77 | + await source.add_schedule(schedule) |
| 78 | + assert await source.get_schedules() == [schedule] |
| 79 | + await source.post_send(schedule) |
| 80 | + assert await source.get_schedules() == [] |
| 81 | + await source.shutdown() |
| 82 | + |
| 83 | + |
| 84 | +@pytest.mark.anyio |
| 85 | +async def test_buffer(redis_url: str) -> None: |
| 86 | + prefix = uuid.uuid4().hex |
| 87 | + source = RedisScheduleSource(redis_url, prefix=prefix, buffer_size=1) |
| 88 | + schedule1 = ScheduledTask( |
| 89 | + task_name="test_task1", |
| 90 | + labels={}, |
| 91 | + args=[], |
| 92 | + kwargs={}, |
| 93 | + cron="* * * * *", |
| 94 | + ) |
| 95 | + schedule2 = ScheduledTask( |
| 96 | + task_name="test_task2", |
| 97 | + labels={}, |
| 98 | + args=[], |
| 99 | + kwargs={}, |
| 100 | + cron="* * * * *", |
| 101 | + ) |
| 102 | + await source.add_schedule(schedule1) |
| 103 | + await source.add_schedule(schedule2) |
| 104 | + schedules = await source.get_schedules() |
| 105 | + assert len(schedules) == 2 |
| 106 | + assert schedule1 in schedules |
| 107 | + assert schedule2 in schedules |
| 108 | + await source.shutdown() |
| 109 | + |
| 110 | + |
| 111 | +@pytest.mark.anyio |
| 112 | +async def test_cluster_set_schedule(redis_cluster_url: str) -> None: |
| 113 | + prefix = uuid.uuid4().hex |
| 114 | + source = RedisClusterScheduleSource(redis_cluster_url, prefix=prefix) |
| 115 | + schedule = ScheduledTask( |
| 116 | + task_name="test_task", |
| 117 | + labels={}, |
| 118 | + args=[], |
| 119 | + kwargs={}, |
| 120 | + cron="* * * * *", |
| 121 | + ) |
| 122 | + await source.add_schedule(schedule) |
| 123 | + schedules = await source.get_schedules() |
| 124 | + assert schedules == [schedule] |
| 125 | + await source.shutdown() |
| 126 | + |
| 127 | + |
| 128 | +@pytest.mark.anyio |
| 129 | +async def test_cluster_delete_schedule(redis_cluster_url: str) -> None: |
| 130 | + prefix = uuid.uuid4().hex |
| 131 | + source = RedisClusterScheduleSource(redis_cluster_url, prefix=prefix) |
| 132 | + schedule = ScheduledTask( |
| 133 | + task_name="test_task", |
| 134 | + labels={}, |
| 135 | + args=[], |
| 136 | + kwargs={}, |
| 137 | + cron="* * * * *", |
| 138 | + ) |
| 139 | + await source.add_schedule(schedule) |
59 | 140 | schedules = await source.get_schedules()
|
60 | 141 | assert schedules == [schedule]
|
| 142 | + await source.delete_schedule(schedule.schedule_id) |
| 143 | + schedules = await source.get_schedules() |
| 144 | + # Schedules are empty. |
| 145 | + assert not schedules |
| 146 | + await source.shutdown() |
| 147 | + |
| 148 | + |
| 149 | +@pytest.mark.anyio |
| 150 | +async def test_cluster_post_run_cron(redis_cluster_url: str) -> None: |
| 151 | + prefix = uuid.uuid4().hex |
| 152 | + source = RedisClusterScheduleSource(redis_cluster_url, prefix=prefix) |
| 153 | + schedule = ScheduledTask( |
| 154 | + task_name="test_task", |
| 155 | + labels={}, |
| 156 | + args=[], |
| 157 | + kwargs={}, |
| 158 | + cron="* * * * *", |
| 159 | + ) |
| 160 | + await source.add_schedule(schedule) |
| 161 | + assert await source.get_schedules() == [schedule] |
| 162 | + await source.post_send(schedule) |
| 163 | + assert await source.get_schedules() == [schedule] |
| 164 | + await source.shutdown() |
| 165 | + |
| 166 | + |
| 167 | +@pytest.mark.anyio |
| 168 | +async def test_cluster_post_run_time(redis_cluster_url: str) -> None: |
| 169 | + prefix = uuid.uuid4().hex |
| 170 | + source = RedisClusterScheduleSource(redis_cluster_url, prefix=prefix) |
| 171 | + schedule = ScheduledTask( |
| 172 | + task_name="test_task", |
| 173 | + labels={}, |
| 174 | + args=[], |
| 175 | + kwargs={}, |
| 176 | + time=dt.datetime(2000, 1, 1), |
| 177 | + ) |
| 178 | + await source.add_schedule(schedule) |
| 179 | + assert await source.get_schedules() == [schedule] |
| 180 | + await source.post_send(schedule) |
| 181 | + assert await source.get_schedules() == [] |
| 182 | + await source.shutdown() |
| 183 | + |
| 184 | + |
| 185 | +@pytest.mark.anyio |
| 186 | +async def test_cluster_buffer(redis_cluster_url: str) -> None: |
| 187 | + prefix = uuid.uuid4().hex |
| 188 | + source = RedisClusterScheduleSource(redis_cluster_url, prefix=prefix, buffer_size=1) |
| 189 | + schedule1 = ScheduledTask( |
| 190 | + task_name="test_task1", |
| 191 | + labels={}, |
| 192 | + args=[], |
| 193 | + kwargs={}, |
| 194 | + cron="* * * * *", |
| 195 | + ) |
| 196 | + schedule2 = ScheduledTask( |
| 197 | + task_name="test_task2", |
| 198 | + labels={}, |
| 199 | + args=[], |
| 200 | + kwargs={}, |
| 201 | + cron="* * * * *", |
| 202 | + ) |
| 203 | + await source.add_schedule(schedule1) |
| 204 | + await source.add_schedule(schedule2) |
| 205 | + schedules = await source.get_schedules() |
| 206 | + assert len(schedules) == 2 |
| 207 | + assert schedule1 in schedules |
| 208 | + assert schedule2 in schedules |
61 | 209 | await source.shutdown()
|
0 commit comments