APP-Backend / Browser_Agent /tests /test_queue.py
Luca448's picture
Savestate: Fix Dockerfile and add Browser_Agent for HF build, bump version to 1.8.334
508316a
Raw History Blame Contribute Delete
4.32 kB
import asyncio
import pytest
from datetime import datetime, timedelta, timezone
from job_queue.task_queue import TaskQueue, TaskStatus
from job_queue.task_store import TaskStore
@pytest.mark.asyncio
async def test_enqueue_returns_task_id():
q = TaskQueue()
task_id = await q.enqueue("agent_task", {"task": "Test"})
assert task_id.startswith("task_")
@pytest.mark.asyncio
async def test_get_status_pending():
q = TaskQueue()
task_id = await q.enqueue("agent_task", {"task": "Test"})
status = await q.get_status(task_id)
assert status["status"] == TaskStatus.PENDING
@pytest.mark.asyncio
async def test_queue_size_limit():
q = TaskQueue(max_size=2)
await q.enqueue("agent_task", {"task": "Test 1"})
await q.enqueue("agent_task", {"task": "Test 2"})
with pytest.raises(asyncio.QueueFull):
await q.enqueue("agent_task", {"task": "Test 3"}, timeout=0.1)
@pytest.mark.asyncio
async def test_cancel_pending_task_updates_public_status():
q = TaskQueue()
task_id = await q.enqueue("agent_task", {"task": "Test"})
assert await q.cancel(task_id) is True
status = await q.get_status(task_id)
assert status["status"] == TaskStatus.CANCELLED
assert status["result"] == "Aufgabe vom Nutzer gestoppt."
assert status["finished_at"] is not None
@pytest.mark.asyncio
async def test_cancel_unknown_task_returns_false():
q = TaskQueue()
assert await q.cancel("task_missing") is False
@pytest.mark.asyncio
async def test_idempotency_key_returns_same_task_and_rejects_payload_change():
q = TaskQueue(store=TaskStore())
first = await q.enqueue(
"agent_task",
{"task": "Research"},
idempotency_key="request-1",
)
duplicate = await q.enqueue(
"agent_task",
{"task": "Research"},
idempotency_key="request-1",
)
assert duplicate == first
with pytest.raises(ValueError, match="idempotency_conflict"):
await q.enqueue(
"agent_task",
{"task": "Different"},
idempotency_key="request-1",
)
@pytest.mark.asyncio
async def test_handoff_resumes_same_task_exactly_once():
q = TaskQueue(store=TaskStore())
task_id = await q.enqueue("agent_task", {"task": "Register"})
task = await q.get_next()
q.task_done()
task.status = TaskStatus.WAITING_FOR_HANDOFF
task.handoff = {
"id": "handoff-1",
"url": "https://example.com/register",
"expires_at": (datetime.now(timezone.utc) + timedelta(minutes=5)).isoformat(),
"completion_policy": "registration_submitted",
"cookies": [],
"form_data": {"email": "agent@example.test"},
}
task.revision = 7
await q.persist(task)
resumed = await q.resume_handoff(
task_id,
"handoff-1",
revision=7,
cookies=[],
resume_url="https://example.com/success",
submission_confirmed=True,
)
duplicate = await q.resume_handoff(
task_id,
"handoff-1",
revision=7,
cookies=[],
resume_url="https://example.com/success",
submission_confirmed=True,
)
assert resumed is duplicate
assert resumed.status == TaskStatus.PENDING
assert resumed.params["_completed_handoff_id"] == "handoff-1"
assert q.size == 1
@pytest.mark.asyncio
async def test_waiting_handoff_survives_queue_reinitialization():
store = TaskStore()
q1 = TaskQueue(store=store)
task_id = await q1.enqueue("agent_task", {"task": "Register"})
task = await q1.get_next()
q1.task_done()
task.status = TaskStatus.WAITING_FOR_HANDOFF
task.handoff = {
"id": "handoff-restart",
"url": "https://example.com/register",
"expires_at": (datetime.now(timezone.utc) + timedelta(minutes=5)).isoformat(),
"completion_policy": "registration_submitted",
"cookies": [{"name": "session", "value": "secret"}],
"form_data": {"email": "agent@example.test", "password": "secret"},
}
await q1.persist(task)
q2 = TaskQueue(store=store)
await q2.initialize()
restored = q2._tasks[task_id]
assert restored.status == TaskStatus.WAITING_FOR_HANDOFF
assert restored.handoff["id"] == "handoff-restart"
assert restored.handoff["form_data"]["password"] == "secret"