-
Notifications
You must be signed in to change notification settings - Fork 1
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Merge pull request #19 from upstash/new-features
Add support for new QStash features
- Loading branch information
Showing
13 changed files
with
724 additions
and
129 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,25 +1,68 @@ | ||
from typing import Callable | ||
|
||
import pytest | ||
|
||
from upstash_qstash import AsyncQStash | ||
|
||
|
||
@pytest.mark.asyncio | ||
async def test_queue_async(async_qstash: AsyncQStash) -> None: | ||
await async_qstash.queue.upsert(queue="test_queue", parallelism=1) | ||
async def test_queue_async( | ||
async_qstash: AsyncQStash, | ||
cleanup_queue_async: Callable[[AsyncQStash, str], None], | ||
) -> None: | ||
name = "test_queue" | ||
cleanup_queue_async(async_qstash, name) | ||
|
||
queue = await async_qstash.queue.get("test_queue") | ||
assert queue.name == "test_queue" | ||
await async_qstash.queue.upsert(queue=name, parallelism=1) | ||
|
||
queue = await async_qstash.queue.get(name) | ||
assert queue.name == name | ||
assert queue.parallelism == 1 | ||
|
||
await async_qstash.queue.upsert(queue="test_queue", parallelism=2) | ||
await async_qstash.queue.upsert(queue=name, parallelism=2) | ||
|
||
queue = await async_qstash.queue.get("test_queue") | ||
assert queue.name == "test_queue" | ||
queue = await async_qstash.queue.get(name) | ||
assert queue.name == name | ||
assert queue.parallelism == 2 | ||
|
||
all_queues = await async_qstash.queue.list() | ||
assert any(True for q in all_queues if q.name == "test_queue") | ||
assert any(True for q in all_queues if q.name == name) | ||
|
||
await async_qstash.queue.delete("test_queue") | ||
await async_qstash.queue.delete(name) | ||
|
||
all_queues = await async_qstash.queue.list() | ||
assert not any(True for q in all_queues if q.name == "test_queue") | ||
assert not any(True for q in all_queues if q.name == name) | ||
|
||
|
||
@pytest.mark.asyncio | ||
async def test_queue_pause_resume_async( | ||
async_qstash: AsyncQStash, | ||
cleanup_queue_async: Callable[[AsyncQStash, str], None], | ||
) -> None: | ||
name = "test_queue" | ||
cleanup_queue_async(async_qstash, name) | ||
|
||
await async_qstash.queue.upsert(queue=name) | ||
|
||
queue = await async_qstash.queue.get(name) | ||
assert queue.paused is False | ||
|
||
await async_qstash.queue.pause(name) | ||
|
||
queue = await async_qstash.queue.get(name) | ||
assert queue.paused is True | ||
|
||
await async_qstash.queue.resume(name) | ||
|
||
queue = await async_qstash.queue.get(name) | ||
assert queue.paused is False | ||
|
||
await async_qstash.queue.upsert(name, paused=True) | ||
|
||
queue = await async_qstash.queue.get(name) | ||
assert queue.paused is True | ||
|
||
await async_qstash.queue.upsert(name, paused=False) | ||
|
||
queue = await async_qstash.queue.get(name) | ||
assert queue.paused is False |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,25 +1,64 @@ | ||
from typing import Callable | ||
|
||
import pytest | ||
|
||
from upstash_qstash import AsyncQStash | ||
|
||
|
||
@pytest.mark.asyncio | ||
async def test_schedule_lifecycle_async(async_qstash: AsyncQStash) -> None: | ||
sched_id = await async_qstash.schedule.create_json( | ||
cron="* * * * *", | ||
async def test_schedule_lifecycle_async( | ||
async_qstash: AsyncQStash, | ||
cleanup_schedule_async: Callable[[AsyncQStash, str], None], | ||
) -> None: | ||
schedule_id = await async_qstash.schedule.create_json( | ||
cron="1 1 1 1 1", | ||
destination="https://example.com", | ||
body={"ex_key": "ex_value"}, | ||
) | ||
|
||
assert len(sched_id) > 0 | ||
cleanup_schedule_async(async_qstash, schedule_id) | ||
|
||
res = await async_qstash.schedule.get(sched_id) | ||
assert res.schedule_id == sched_id | ||
assert res.cron == "* * * * *" | ||
assert len(schedule_id) > 0 | ||
|
||
res = await async_qstash.schedule.get(schedule_id) | ||
assert res.schedule_id == schedule_id | ||
assert res.cron == "1 1 1 1 1" | ||
|
||
list_res = await async_qstash.schedule.list() | ||
assert any(s.schedule_id == sched_id for s in list_res) | ||
assert any(s.schedule_id == schedule_id for s in list_res) | ||
|
||
await async_qstash.schedule.delete(sched_id) | ||
await async_qstash.schedule.delete(schedule_id) | ||
|
||
list_res = await async_qstash.schedule.list() | ||
assert not any(s.schedule_id == sched_id for s in list_res) | ||
assert not any(s.schedule_id == schedule_id for s in list_res) | ||
|
||
|
||
@pytest.mark.asyncio | ||
async def test_schedule_pause_resume_async( | ||
async_qstash: AsyncQStash, | ||
cleanup_schedule_async: Callable[[AsyncQStash, str], None], | ||
) -> None: | ||
schedule_id = await async_qstash.schedule.create_json( | ||
cron="1 1 1 1 1", | ||
destination="https://example.com", | ||
body={"ex_key": "ex_value"}, | ||
) | ||
|
||
cleanup_schedule_async(async_qstash, schedule_id) | ||
|
||
assert len(schedule_id) > 0 | ||
|
||
res = await async_qstash.schedule.get(schedule_id) | ||
assert res.schedule_id == schedule_id | ||
assert res.cron == "1 1 1 1 1" | ||
assert res.paused is False | ||
|
||
await async_qstash.schedule.pause(schedule_id) | ||
|
||
res = await async_qstash.schedule.get(schedule_id) | ||
assert res.paused is True | ||
|
||
await async_qstash.schedule.resume(schedule_id) | ||
|
||
res = await async_qstash.schedule.get(schedule_id) | ||
assert res.paused is False |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.