Coverage for streamwise/node_manager.py: 78%

60 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-09 04:47 +0000

1""" 

2Module to manage Kubernetes nodes. 

3""" 

4 

5import sys 

6import logging 

7import traceback 

8 

9from http import HTTPStatus 

10 

11from quart import jsonify 

12from quart import abort 

13from quart import render_template 

14 

15from typing import Optional 

16 

17from kubernetes_asyncio.client import ApiClient 

18from kubernetes_asyncio.client import CoreV1Api 

19from kubernetes_asyncio.client import V1DeleteOptions 

20 

21sys.path.append("..") 

22from quart_utils import QuartReturn 

23 

24from k8s_utils import load_k8s_config 

25from k8s_utils import get_k8s_nodes 

26from k8s_utils import get_k8s_pods 

27 

28 

29async def remove_node( 

30 node_name: str, 

31 k8s_cluster: Optional[str] = None 

32) -> QuartReturn: 

33 """API interface to remove a node by name.""" 

34 if not node_name: 

35 return jsonify({"error": "Node name is required"}), HTTPStatus.BAD_REQUEST 

36 

37 await load_k8s_config(k8s_cluster) 

38 async with ApiClient() as api_client: 

39 k8s_api = CoreV1Api(api_client) 

40 try: 

41 await k8s_api.delete_node(name=node_name, body=V1DeleteOptions()) 

42 pods = await k8s_api.list_pod_for_all_namespaces(field_selector=f"spec.nodeName={node_name}") 

43 for pod in pods.items: 

44 await k8s_api.delete_namespaced_pod( 

45 name=pod.metadata.name, 

46 namespace=pod.metadata.namespace) 

47 return jsonify({"message": f"Node {node_name} removed successfully"}), HTTPStatus.OK 

48 except Exception as ex: 

49 logging.error(f"Error removing node {node_name}: {ex}.") 

50 return jsonify({"error": str(ex)}), HTTPStatus.INTERNAL_SERVER_ERROR 

51 

52 

53async def node_info( 

54 node_name: str, 

55 k8s_cluster: Optional[str] = None 

56) -> QuartReturn: 

57 """Display information about a specific node.""" 

58 if not node_name: 

59 return jsonify({"error": "Node name is required"}), HTTPStatus.BAD_REQUEST 

60 try: 

61 nodes = await get_k8s_nodes(k8s_cluster) 

62 if not nodes: 

63 return jsonify({"error": f"Node '{node_name}' not found"}), HTTPStatus.NOT_FOUND 

64 for node in nodes: 

65 if node["node_name"] == node_name: 

66 pods = await get_k8s_pods(k8s_cluster) 

67 node["pods"] = [pod for pod in pods if pod["node"] == node_name] 

68 return await render_template( 

69 "node.html", 

70 nodes=[node]) 

71 except Exception as ex: 

72 logging.error(f"Error fetching node info for {node_name}: {ex}.") 

73 abort( 

74 HTTPStatus.INTERNAL_SERVER_ERROR, 

75 description=f"Error fetching node info for {node_name}: {ex}") 

76 return jsonify({"error": f"Node '{node_name}' not found"}), HTTPStatus.NOT_FOUND 

77 

78 

79async def nodes_info( 

80 k8s_cluster: Optional[str] = None 

81) -> QuartReturn: 

82 """Display information about all nodes.""" 

83 try: 

84 nodes = await get_k8s_nodes(k8s_cluster) 

85 if not nodes: 

86 return jsonify({"error": "No nodes found"}), HTTPStatus.NOT_FOUND 

87 

88 pods = await get_k8s_pods(k8s_cluster) 

89 for node in nodes: 

90 node_name = node["node_name"] 

91 node["pods"] = [pod for pod in pods if pod["node"] == node_name] 

92 return await render_template( 

93 "node.html", 

94 nodes=nodes) 

95 except Exception as ex: 

96 logging.error(f"Error fetching nodes info [{type(ex)}]: {ex}.") 

97 return jsonify({ 

98 "error": str(ex), 

99 "type": str(type(ex)), 

100 "trace": traceback.format_exc(), 

101 }), HTTPStatus.INTERNAL_SERVER_ERROR