File size: 4,319 Bytes
508316a
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
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"