Coverage for tests/streamwise_app/test_client.py: 100%
61 statements
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-09 04:47 +0000
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-09 04:47 +0000
1#!/usr/bin/env python3
3import os
4import sys
5import pytest
6import asyncio
8from unittest.mock import patch
10# Add current path
11sys.path.append(os.getcwd())
13from tests.test_utils import temp_sys_path
15mock_modules: dict[str, object] = {
16 # "torch": mock_torch,
17}
19from k8s_utils import K8sService
20from k8s_utils import K8sContainer
22with patch.dict(sys.modules, mock_modules):
23 with temp_sys_path("apps", "apps/streampersona"):
24 from apps.client import ServiceRequest
25 from apps.client import ServiceRequestWorker
26 from apps.client import ServiceError
27 from apps.lmm_service_manager import LMMServiceManager
28 from k8s_utils import ServiceNotFoundError
31def test_service_request() -> None:
32 req = ServiceRequest(
33 service_name="test_service",
34 request_id="req_123",
35 payload_json={"key": "value"},
36 url="http://localhost:8000/test_service"
37 )
38 assert req.service_name == "test_service"
40 assert req.is_running() is False
41 assert req.client_timeout.connect is None
42 assert req.get_base_request_url() == "http://localhost:8000"
43 assert req.model_dump()["url"] == "http://localhost:8000/test_service"
44 req_json = req.json()
45 assert req_json is not None
46 assert req_json.startswith('{')
47 assert req_json.endswith('}')
48 assert req.get_payload_len() > 0
50 # req2 = ServiceRequest.parse_json(req_json)
53@pytest.mark.asyncio
54async def test_service_request_worker() -> None:
55 service_manager = LMMServiceManager("streamwise")
56 worker = ServiceRequestWorker(
57 app_name="TestApp",
58 service_manager=service_manager,
59 )
60 try:
61 # Skip for unit test
62 # await worker.start()
64 request = ServiceRequest(
65 service_name="test_service",
66 request_id="req_123",
67 payload_json={"key": "value"},
68 url="http://localhost:8000/test_service"
69 )
71 with pytest.raises(ServiceNotFoundError):
72 await worker.submit_request(request)
74 service = K8sService("test_service")
75 service_manager.services["test_service"] = service
76 future = await worker.submit_request(request)
77 assert future is not None
78 assert isinstance(future, asyncio.Future)
79 assert request.future is not None and request.future.done() is False
80 assert request.status == "CREATED"
81 assert request.exception is None
83 await worker._http_request(request)
84 assert request.future is not None and request.future.done() is True
85 assert request.status == "FAILED"
86 assert request.exception is not None
87 assert isinstance(request.exception, ServiceError)
88 assert "No active containers for test_service" in str(request.exception)
89 assert request.retries == 0
91 # Success request
92 request_1 = ServiceRequest(
93 service_name="test_service",
94 request_id="req_124",
95 payload_json={"key": "value"},
96 url="http://localhost:8000/test_service"
97 )
99 container = K8sContainer("test_service_1234")
100 service.add_container(container)
102 await worker._http_request(request_1)
103 assert request_1.future is None
104 assert request_1.status == "RETRYING"
105 assert request_1.exception is None
106 assert request_1.retries == 1
107 finally:
108 await worker.stop()