Spaces:
Paused
Paused
| 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 | |