File size: 2,276 Bytes
d958e80
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
import pytest

from app.job_queue import create_job, assign_job_to_worker, complete_job, fail_job, get_job, select_worker_for_job
from app.models import JobType, PrivacyMode, JobStatus, WorkerState, WorkerRuntimeType, JobResult
from app.session_store import create_session, attach_worker
from app.security import generate_worker_id


def test_create_job():
    session = create_session()
    job = create_job(
        session_id=session.session_id,
        job_type=JobType.TEXT_EMBEDDING,
        payload={"text": "hello"},
        privacy_mode=PrivacyMode.RAW_INPUT_REMOTE,
        constraints={},
    )
    assert job.job_id.startswith("job_")
    assert job.status == JobStatus.QUEUED


def test_assign_job_to_worker():
    session = create_session()
    job = create_job(session.session_id, JobType.TEXT_EMBEDDING, {"text": "x"}, PrivacyMode.RAW_INPUT_REMOTE, {})
    ok = assign_job_to_worker(job.job_id, "wk_123")
    assert ok is True
    updated = get_job(job.job_id)
    assert updated.worker_id == "wk_123"
    assert updated.status == JobStatus.ASSIGNED


def test_complete_job():
    session = create_session()
    job = create_job(session.session_id, JobType.TEXT_EMBEDDING, {"text": "x"}, PrivacyMode.RAW_INPUT_REMOTE, {})
    assign_job_to_worker(job.job_id, "wk_123")
    result = JobResult(job_id=job.job_id, worker_id="wk_123", output={"vec": [0.1]}, latency_ms=100)
    ok = complete_job(job.job_id, result)
    assert ok is True
    updated = get_job(job.job_id)
    assert updated.status == JobStatus.COMPLETED
    assert updated.result == {"vec": [0.1]}


def test_fail_job():
    session = create_session()
    job = create_job(session.session_id, JobType.TEXT_EMBEDDING, {"text": "x"}, PrivacyMode.RAW_INPUT_REMOTE, {})
    ok = fail_job(job.job_id, "worker crashed")
    assert ok is True
    updated = get_job(job.job_id)
    assert updated.status == JobStatus.FAILED
    assert updated.error_reason == "worker crashed"


def test_reject_job_when_no_worker():
    session = create_session()
    job = create_job(session.session_id, JobType.TEXT_EMBEDDING, {"text": "x"}, PrivacyMode.LOCAL_ONLY_RESULT_ONLY, {})
    worker_id = select_worker_for_job(session.session_id, job)
    # No workers attached, should return None
    assert worker_id is None