Skip to content

Commit 190acbe

Browse files
committed
Resolve #25 -- Add simple task/queue actions to the inspector
1 parent f9749f8 commit 190acbe

9 files changed

Lines changed: 971 additions & 98 deletions

File tree

‎tests/backends/test_base.py‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -86,6 +86,16 @@ def test_telemetry__raise_not_implemented_error(self) -> None:
8686
with pytest.raises(NotImplementedError):
8787
BackendDouble(alias="default", params={}).telemetry()
8888

89+
def test_dequeue__raise_not_implemented_error(self) -> None:
90+
"""Raise NotImplementedError for backend dequeue API."""
91+
with pytest.raises(NotImplementedError):
92+
BackendDouble(alias="default", params={}).dequeue(task_result=None)
93+
94+
def test_purge__raise_not_implemented_error(self) -> None:
95+
"""Raise NotImplementedError for backend purge_queue API."""
96+
with pytest.raises(NotImplementedError):
97+
BackendDouble(alias="default", params={}).purge("default")
98+
8999
def test_validate_task__accepts_module_level_function(self) -> None:
90100
"""validate_task accepts a module-level retry callback."""
91101
backend = BackendDouble(alias="default", params={})

‎tests/backends/test_redis.py‎

Lines changed: 230 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -89,12 +89,10 @@ def test_main__trims_stale_telemetry(self):
8989
ingress_key = backend.INGRESS_KEY.format(
9090
prefix=backend.key_prefix, queue_name="default"
9191
)
92-
successful_key = backend.SUCCESSFUL_RESULTS_KEY.format(
93-
prefix=backend.key_prefix, queue_name="default"
94-
)
95-
failed_key = backend.FAILED_RESULTS_KEY.format(
96-
prefix=backend.key_prefix, queue_name="default"
92+
successful_key = backend._segment_key(
93+
TaskResultStatus.SUCCESSFUL, "default"
9794
)
95+
failed_key = backend._segment_key(TaskResultStatus.FAILED, "default")
9896
old = (timezone.now() - datetime.timedelta(seconds=120)).timestamp() * 1000
9997
backend.client.zadd(ingress_key, {task_result.id: old})
10098
backend.client.zadd(successful_key, {task_result.id: old})
@@ -138,9 +136,7 @@ def test_acquire__moves_to_running_set(self):
138136
assert acquired.id == task_result.id
139137

140138
# Verify task is in running set, not in any processing set
141-
running_key = backend.RUNNING_KEY.format(
142-
prefix=backend.key_prefix, queue_name="default"
143-
)
139+
running_key = backend._segment_key(TaskResultStatus.RUNNING, "default")
144140
assert backend.client.zscore(running_key, task_result.id) is not None
145141

146142
# Verify task data was updated with worker info
@@ -401,12 +397,10 @@ def test_telemetry__successful_failed_evicted_by_result_ttl(self):
401397
},
402398
)
403399
try:
404-
successful_key = backend.SUCCESSFUL_RESULTS_KEY.format(
405-
prefix=backend.key_prefix, queue_name="default"
406-
)
407-
failed_key = backend.FAILED_RESULTS_KEY.format(
408-
prefix=backend.key_prefix, queue_name="default"
400+
successful_key = backend._segment_key(
401+
TaskResultStatus.SUCCESSFUL, "default"
409402
)
403+
failed_key = backend._segment_key(TaskResultStatus.FAILED, "default")
410404

411405
def _ack(status: TaskResultStatus) -> str:
412406
enqueued = backend.enqueue(echo, args=[1])
@@ -472,8 +466,8 @@ def test_telemetry__ingress_egress_age_out_of_window(self):
472466
ingress_key = backend.INGRESS_KEY.format(
473467
prefix=backend.key_prefix, queue_name="default"
474468
)
475-
successful_results_key = backend.SUCCESSFUL_RESULTS_KEY.format(
476-
prefix=backend.key_prefix, queue_name="default"
469+
successful_results_key = backend._segment_key(
470+
TaskResultStatus.SUCCESSFUL, "default"
477471
)
478472
old = (timezone.now() - datetime.timedelta(seconds=120)).timestamp() * 1000
479473
backend.client.zadd(ingress_key, {task_result.id: old})
@@ -556,8 +550,8 @@ def test_telemetry__egress_window_follows_display_interval(self):
556550
finished_at=timezone.now(),
557551
)
558552
)
559-
successful_key = backend.SUCCESSFUL_RESULTS_KEY.format(
560-
prefix=backend.key_prefix, queue_name="default"
553+
successful_key = backend._segment_key(
554+
TaskResultStatus.SUCCESSFUL, "default"
561555
)
562556

563557
# Inside result_ttl (60s) but outside a 10s display window.
@@ -720,9 +714,7 @@ def test_requeue__moves_from_running_to_deferred(self) -> None:
720714
run_after = timezone.now() + datetime.timedelta(seconds=10)
721715
backend.requeue(failed, run_after)
722716

723-
running_key = backend.RUNNING_KEY.format(
724-
prefix=backend.key_prefix, queue_name="default"
725-
)
717+
running_key = backend._segment_key(TaskResultStatus.RUNNING, "default")
726718
deferred_key = backend.DEFERRED_KEY.format(
727719
prefix=backend.key_prefix, queue_name="default"
728720
)
@@ -832,3 +824,221 @@ def test_requeue__task_is_re_acquirable_after_delay(self) -> None:
832824
assert re_acquired.id == task_result.id
833825
finally:
834826
backend.close()
827+
828+
def test_requeue__cleans_up_failed_and_result_keys(self) -> None:
829+
"""requeue() removes the task from the failed zset and deletes its result key."""
830+
backend = RedisTaskBackend(
831+
"requeue_cleanup_test",
832+
{
833+
"QUEUES": ["default"],
834+
"REDIS_URL": "redis://localhost:6379/0",
835+
"OPTIONS": {
836+
"lease_ttl": datetime.timedelta(hours=1),
837+
"result_ttl": datetime.timedelta(seconds=60),
838+
},
839+
},
840+
)
841+
try:
842+
backend.enqueue(echo, args=[1])
843+
acquired = backend.acquire(
844+
timeout=datetime.timedelta(seconds=1), worker="cleanup-test"
845+
)
846+
assert acquired is not None
847+
backend.acknowledge(
848+
dataclasses.replace(
849+
acquired,
850+
status=TaskResultStatus.FAILED,
851+
finished_at=timezone.now(),
852+
)
853+
)
854+
failed = next(
855+
backend.peek(
856+
queue_name="default", status=TaskResultStatus.FAILED, count=10
857+
)
858+
)
859+
failed_key = backend._segment_key(TaskResultStatus.FAILED, "default")
860+
result_key = backend.RESULT_KEY.format(
861+
prefix=backend.key_prefix, result_id=failed.id
862+
)
863+
assert backend.client.zscore(failed_key, failed.id) is not None
864+
assert backend.client.exists(result_key)
865+
866+
backend.requeue(failed, timezone.now() + datetime.timedelta(seconds=10))
867+
868+
assert backend.client.zscore(failed_key, failed.id) is None
869+
assert not backend.client.exists(result_key)
870+
finally:
871+
backend.close()
872+
873+
def test_dequeue__removes_ready_task_from_queue(self) -> None:
874+
"""dequeue() removes a ready task from the queue zset."""
875+
backend = RedisTaskBackend(
876+
"dequeue_ready_test",
877+
{
878+
"QUEUES": ["default"],
879+
"REDIS_URL": "redis://localhost:6379/0",
880+
"OPTIONS": {
881+
"lease_ttl": datetime.timedelta(hours=1),
882+
"result_ttl": datetime.timedelta(seconds=60),
883+
},
884+
},
885+
)
886+
try:
887+
task_result = backend.enqueue(echo, args=[1])
888+
queue_key = backend._segment_key(TaskResultStatus.READY, "default")
889+
assert backend.client.zscore(queue_key, task_result.id) is not None
890+
891+
backend.dequeue(task_result)
892+
893+
assert backend.client.zscore(queue_key, task_result.id) is None
894+
finally:
895+
backend.close()
896+
897+
def test_dequeue__removes_failed_task_from_results(self) -> None:
898+
"""dequeue() removes a failed task from the failed zset."""
899+
backend = RedisTaskBackend(
900+
"dequeue_failed_test",
901+
{
902+
"QUEUES": ["default"],
903+
"REDIS_URL": "redis://localhost:6379/0",
904+
"OPTIONS": {
905+
"lease_ttl": datetime.timedelta(hours=1),
906+
"result_ttl": datetime.timedelta(seconds=60),
907+
},
908+
},
909+
)
910+
try:
911+
backend.enqueue(echo, args=[1])
912+
acquired = backend.acquire(
913+
timeout=datetime.timedelta(seconds=1), worker="dequeue-failed-test"
914+
)
915+
assert acquired is not None
916+
backend.acknowledge(
917+
dataclasses.replace(
918+
acquired,
919+
status=TaskResultStatus.FAILED,
920+
finished_at=timezone.now(),
921+
)
922+
)
923+
failed = next(
924+
backend.peek(
925+
queue_name="default", status=TaskResultStatus.FAILED, count=10
926+
)
927+
)
928+
failed_key = backend._segment_key(TaskResultStatus.FAILED, "default")
929+
930+
backend.dequeue(failed)
931+
932+
assert backend.client.zscore(failed_key, failed.id) is None
933+
finally:
934+
backend.close()
935+
936+
def test_purge__removes_all_tasks_across_segments(self) -> None:
937+
"""purge_queue() deletes every task across all segments."""
938+
backend = RedisTaskBackend(
939+
"purge_test",
940+
{
941+
"QUEUES": ["default"],
942+
"REDIS_URL": "redis://localhost:6379/0",
943+
"OPTIONS": {
944+
"lease_ttl": datetime.timedelta(hours=1),
945+
"result_ttl": datetime.timedelta(seconds=60),
946+
},
947+
},
948+
)
949+
try:
950+
# Two ready tasks
951+
backend.enqueue(echo, args=[1])
952+
backend.enqueue(echo, args=[2])
953+
# One running task
954+
backend.acquire(timeout=datetime.timedelta(seconds=1), worker="purge-test")
955+
# One failed task
956+
backend.enqueue(echo, args=[3])
957+
acquired = backend.acquire(
958+
timeout=datetime.timedelta(seconds=1), worker="purge-test-2"
959+
)
960+
assert acquired is not None
961+
backend.acknowledge(
962+
dataclasses.replace(
963+
acquired,
964+
status=TaskResultStatus.FAILED,
965+
finished_at=timezone.now(),
966+
)
967+
)
968+
# One successful task
969+
backend.enqueue(echo, args=[4])
970+
acquired = backend.acquire(
971+
timeout=datetime.timedelta(seconds=1), worker="purge-test-3"
972+
)
973+
assert acquired is not None
974+
backend.acknowledge(
975+
dataclasses.replace(
976+
acquired,
977+
status=TaskResultStatus.SUCCESSFUL,
978+
finished_at=timezone.now(),
979+
)
980+
)
981+
982+
backend.purge("default")
983+
assert (
984+
list(
985+
backend.peek(
986+
queue_name="default", status=TaskResultStatus.READY, count=10
987+
)
988+
)
989+
== []
990+
)
991+
assert (
992+
list(
993+
backend.peek(
994+
queue_name="default", status=TaskResultStatus.RUNNING, count=10
995+
)
996+
)
997+
== []
998+
)
999+
assert (
1000+
list(
1001+
backend.peek(
1002+
queue_name="default", status=TaskResultStatus.FAILED, count=10
1003+
)
1004+
)
1005+
== []
1006+
)
1007+
assert (
1008+
list(
1009+
backend.peek(
1010+
queue_name="default",
1011+
status=TaskResultStatus.SUCCESSFUL,
1012+
count=10,
1013+
)
1014+
)
1015+
== []
1016+
)
1017+
finally:
1018+
backend.close()
1019+
1020+
def test_purge__empty_queue_is_noop(self) -> None:
1021+
"""purge_queue() on an empty queue is a no-op."""
1022+
backend = RedisTaskBackend(
1023+
"purge_empty_test",
1024+
{
1025+
"QUEUES": ["default"],
1026+
"REDIS_URL": "redis://localhost:6379/0",
1027+
"OPTIONS": {
1028+
"lease_ttl": datetime.timedelta(hours=1),
1029+
"result_ttl": datetime.timedelta(seconds=60),
1030+
},
1031+
},
1032+
)
1033+
try:
1034+
backend.purge("default")
1035+
assert (
1036+
list(
1037+
backend.peek(
1038+
queue_name="default", status=TaskResultStatus.READY, count=10
1039+
)
1040+
)
1041+
== []
1042+
)
1043+
finally:
1044+
backend.close()

0 commit comments

Comments
 (0)