Coverage for streamwise/service_account_manager.py: 73%

56 statements  

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

1""" 

2Service account manager for Kubernetes clusters. 

3""" 

4 

5import sys 

6 

7from http import HTTPStatus 

8 

9from typing import List 

10from typing import Optional 

11 

12from kubernetes_asyncio.client import ApiClient 

13from kubernetes_asyncio.client import ApiException 

14from kubernetes_asyncio.client import CoreV1Api 

15from kubernetes_asyncio.client import V1ObjectMeta 

16from kubernetes_asyncio.client import V1PolicyRule 

17from kubernetes_asyncio.client import V1ClusterRole 

18from kubernetes_asyncio.client import V1ClusterRoleBinding 

19from kubernetes_asyncio.client import RbacAuthorizationV1Api 

20from kubernetes_asyncio.client import V1ServiceAccount 

21from kubernetes_asyncio.client import V1RoleRef 

22from kubernetes_asyncio.client import RbacV1Subject 

23 

24sys.path.append("..") 

25from k8s_utils import load_k8s_config 

26 

27 

28async def ensure_service_account( 

29 core_api: CoreV1Api, 

30 service_account_name: str, 

31 namespace: str 

32) -> None: 

33 try: 

34 await core_api.read_namespaced_service_account(service_account_name, namespace) 

35 except ApiException as ex: 

36 if ex.status != HTTPStatus.NOT_FOUND: 

37 raise 

38 sa_body = V1ServiceAccount( 

39 metadata=V1ObjectMeta(name=service_account_name)) 

40 await core_api.create_namespaced_service_account(namespace=namespace, body=sa_body) 

41 

42 

43async def ensure_cluster_role( 

44 rbac_api: RbacAuthorizationV1Api, 

45 name: str, 

46 rules: List[V1PolicyRule] 

47) -> None: 

48 try: 

49 await rbac_api.read_cluster_role(name) 

50 except ApiException as ex: 

51 if ex.status != HTTPStatus.NOT_FOUND: 

52 raise 

53 cluster_role_body = V1ClusterRole( 

54 metadata=V1ObjectMeta(name=name), 

55 rules=rules) 

56 await rbac_api.create_cluster_role(body=cluster_role_body) 

57 

58 

59async def ensure_cluster_role_binding( 

60 rbac_api: RbacAuthorizationV1Api, 

61 binding_name: str, 

62 role_name: str, 

63 service_account_name: str, 

64 namespace: str 

65) -> None: 

66 try: 

67 await rbac_api.read_cluster_role_binding(binding_name) 

68 except ApiException as ex: 

69 if ex.status != HTTPStatus.NOT_FOUND: 

70 raise 

71 crb_body = V1ClusterRoleBinding( 

72 metadata=V1ObjectMeta(name=binding_name), 

73 role_ref=V1RoleRef( 

74 kind="ClusterRole", 

75 name=role_name, 

76 api_group="rbac.authorization.k8s.io", 

77 ), 

78 subjects=[RbacV1Subject( 

79 kind="ServiceAccount", 

80 name=service_account_name, 

81 namespace=namespace, 

82 )] 

83 ) 

84 await rbac_api.create_cluster_role_binding(body=crb_body) 

85 

86 

87async def get_service_account( 

88 k8s_cluster: Optional[str], 

89 service_account_name: str, 

90 cluster_role_name: str, 

91 cluster_role_binding_name: str, 

92 rules: List[V1PolicyRule], 

93 namespace: str = "default" 

94) -> str: 

95 """ 

96 Create a service account with specified permissions if it doesn't exist. 

97 """ 

98 await load_k8s_config(k8s_cluster) 

99 async with ApiClient() as api_client: 

100 core_api = CoreV1Api(api_client) 

101 await ensure_service_account( 

102 core_api, service_account_name, namespace) 

103 

104 rbac_api = RbacAuthorizationV1Api(api_client) 

105 await ensure_cluster_role( 

106 rbac_api, cluster_role_name, rules) 

107 await ensure_cluster_role_binding( 

108 rbac_api, cluster_role_binding_name, cluster_role_name, service_account_name, namespace) 

109 return service_account_name 

110 

111 

112async def get_streamwiseapp_service_account( 

113 k8s_cluster: Optional[str] = None, 

114 namespace: str = "default" 

115) -> str: 

116 """ 

117 Create a service account with necessary permissions for Streamcast if it doesn't exist. 

118 This allow listing nodes, pods, services, events, and namespaces. 

119 """ 

120 rules = [ 

121 V1PolicyRule( 

122 api_groups=[""], 

123 resources=[ 

124 "nodes", 

125 "pods", 

126 "pods/log", 

127 "services", 

128 "events", 

129 "namespaces", 

130 "clusterroles", 

131 ], 

132 verbs=["get", "list", "watch"] 

133 ) 

134 ] 

135 return await get_service_account( 

136 k8s_cluster=k8s_cluster, 

137 service_account_name="streamwiseapp-service-account", 

138 cluster_role_name="streamwiseapp-manager", 

139 cluster_role_binding_name="streamwiseapp-manager-binding", 

140 rules=rules, 

141 namespace=namespace) 

142 

143 

144async def get_streamwise_service_account( 

145 k8s_cluster: Optional[str] = None, 

146 namespace: str = "default" 

147) -> str: 

148 """ 

149 Create a service account with necessary permissions for Streamwise if it doesn't exist. 

150 This allow listing and creating namespaces, services, pods, and nodes. 

151 """ 

152 rules = [ 

153 V1PolicyRule( 

154 api_groups=[""], 

155 resources=[ 

156 "nodes", 

157 "pods", 

158 "pods/log", 

159 "services", 

160 "events", 

161 "namespaces", 

162 "serviceaccounts", 

163 "clusterroles", 

164 "clusterrolebindings", 

165 ], 

166 verbs=["create", "get", "list", "watch", "delete", "patch"] 

167 ) 

168 ] 

169 return await get_service_account( 

170 k8s_cluster=k8s_cluster, 

171 service_account_name="streamwise-service-account", 

172 cluster_role_name="streamwise-manager", 

173 cluster_role_binding_name="streamwise-manager-binding", 

174 rules=rules, 

175 namespace=namespace)