|
1 | 1 | import logging |
| 2 | +from datetime import datetime, timedelta, timezone |
2 | 3 |
|
3 | 4 | import pytest |
4 | 5 |
|
5 | 6 | from docket import CurrentDocket, CurrentWorker, Docket, Worker |
6 | | -from docket.dependencies import Depends, Retry, TaskArgument |
| 7 | +from docket.dependencies import Depends, ExponentialRetry, Retry, TaskArgument |
7 | 8 |
|
8 | 9 |
|
9 | 10 | async def test_dependencies_may_be_duplicated(docket: Docket, worker: Worker): |
@@ -95,6 +96,127 @@ async def the_task( |
95 | 96 | assert calls == 2 |
96 | 97 |
|
97 | 98 |
|
| 99 | +@pytest.mark.parametrize("retry_cls", [Retry, ExponentialRetry]) |
| 100 | +async def test_user_can_request_a_retry_in_timedelta_time( |
| 101 | + retry_cls: Retry, docket: Docket, worker: Worker |
| 102 | +): |
| 103 | + calls = 0 |
| 104 | + first_call_time = None |
| 105 | + second_call_time = None |
| 106 | + |
| 107 | + async def the_task( |
| 108 | + a: str, |
| 109 | + b: str, |
| 110 | + retry: Retry = retry_cls(attempts=2), # type: ignore[reportCallIssue] |
| 111 | + ): |
| 112 | + assert a == "a" |
| 113 | + assert b == "b" |
| 114 | + |
| 115 | + nonlocal calls |
| 116 | + calls += 1 |
| 117 | + |
| 118 | + nonlocal first_call_time |
| 119 | + if not first_call_time: |
| 120 | + first_call_time = datetime.now(timezone.utc) |
| 121 | + retry.in_(timedelta(seconds=0.5)) |
| 122 | + else: |
| 123 | + nonlocal second_call_time |
| 124 | + second_call_time = datetime.now(timezone.utc) |
| 125 | + |
| 126 | + await docket.add(the_task)("a", "b") |
| 127 | + |
| 128 | + await worker.run_until_finished() |
| 129 | + |
| 130 | + assert calls == 2 |
| 131 | + |
| 132 | + assert isinstance(first_call_time, datetime) |
| 133 | + assert isinstance(second_call_time, datetime) |
| 134 | + |
| 135 | + delay = second_call_time - first_call_time |
| 136 | + assert delay.total_seconds() > 0 < 1 |
| 137 | + |
| 138 | + |
| 139 | +@pytest.mark.parametrize("retry_cls", [Retry, ExponentialRetry]) |
| 140 | +async def test_user_can_request_a_retry_at_a_specific_time( |
| 141 | + retry_cls: Retry, docket: Docket, worker: Worker |
| 142 | +): |
| 143 | + calls = 0 |
| 144 | + first_call_time = None |
| 145 | + second_call_time = None |
| 146 | + |
| 147 | + async def the_task( |
| 148 | + a: str, |
| 149 | + b: str, |
| 150 | + retry: Retry = retry_cls(attempts=2), # type: ignore[reportCallIssue] |
| 151 | + ): |
| 152 | + assert a == "a" |
| 153 | + assert b == "b" |
| 154 | + |
| 155 | + nonlocal calls |
| 156 | + calls += 1 |
| 157 | + |
| 158 | + nonlocal first_call_time |
| 159 | + if not first_call_time: |
| 160 | + when = datetime.now(timezone.utc) + timedelta(seconds=0.5) |
| 161 | + first_call_time = datetime.now(timezone.utc) |
| 162 | + retry.at(when) |
| 163 | + else: |
| 164 | + nonlocal second_call_time |
| 165 | + second_call_time = datetime.now(timezone.utc) |
| 166 | + |
| 167 | + await docket.add(the_task)("a", "b") |
| 168 | + |
| 169 | + await worker.run_until_finished() |
| 170 | + |
| 171 | + assert calls == 2 |
| 172 | + |
| 173 | + assert isinstance(first_call_time, datetime) |
| 174 | + assert isinstance(second_call_time, datetime) |
| 175 | + |
| 176 | + delay = second_call_time - first_call_time |
| 177 | + assert delay.total_seconds() > 0 < 1 |
| 178 | + |
| 179 | + |
| 180 | +async def test_user_can_request_a_retry_at_a_specific_time_in_the_past( |
| 181 | + docket: Docket, worker: Worker |
| 182 | +): |
| 183 | + calls = 0 |
| 184 | + first_call_time = None |
| 185 | + second_call_time = None |
| 186 | + |
| 187 | + async def the_task( |
| 188 | + a: str, |
| 189 | + b: str, |
| 190 | + retry: Retry = Retry(attempts=2), |
| 191 | + ): |
| 192 | + assert a == "a" |
| 193 | + assert b == "b" |
| 194 | + |
| 195 | + nonlocal calls |
| 196 | + calls += 1 |
| 197 | + |
| 198 | + nonlocal first_call_time |
| 199 | + if not first_call_time: |
| 200 | + when = datetime.now(timezone.utc) - timedelta(days=1) |
| 201 | + first_call_time = datetime.now(timezone.utc) |
| 202 | + retry.at(when) |
| 203 | + else: |
| 204 | + nonlocal second_call_time |
| 205 | + second_call_time = datetime.now(timezone.utc) |
| 206 | + |
| 207 | + await docket.add(the_task)("a", "b") |
| 208 | + |
| 209 | + await worker.run_until_finished() |
| 210 | + |
| 211 | + assert calls == 2 |
| 212 | + |
| 213 | + assert isinstance(first_call_time, datetime) |
| 214 | + assert isinstance(second_call_time, datetime) |
| 215 | + |
| 216 | + delay = second_call_time - first_call_time |
| 217 | + assert delay.total_seconds() > 0 < 1 |
| 218 | + |
| 219 | + |
98 | 220 | async def test_dependencies_error_for_missing_task_argument( |
99 | 221 | docket: Docket, worker: Worker, caplog: pytest.LogCaptureFixture |
100 | 222 | ): |
|
0 commit comments