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"