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

1#!/usr/bin/env python3 

2 

3import os 

4import sys 

5import pytest 

6import asyncio 

7 

8from unittest.mock import patch 

9 

10# Add current path 

11sys.path.append(os.getcwd()) 

12 

13from tests.test_utils import temp_sys_path 

14 

15mock_modules: dict[str, object] = { 

16 # "torch": mock_torch, 

17} 

18 

19from k8s_utils import K8sService 

20from k8s_utils import K8sContainer 

21 

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 

29 

30 

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" 

39 

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 

49 

50 # req2 = ServiceRequest.parse_json(req_json) 

51 

52 

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() 

63 

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 ) 

70 

71 with pytest.raises(ServiceNotFoundError): 

72 await worker.submit_request(request) 

73 

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 

82 

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 

90 

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 ) 

98 

99 container = K8sContainer("test_service_1234") 

100 service.add_container(container) 

101 

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()