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
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-09 04:47 +0000
1"""
2Service account manager for Kubernetes clusters.
3"""
5import sys
7from http import HTTPStatus
9from typing import List
10from typing import Optional
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
24sys.path.append("..")
25from k8s_utils import load_k8s_config
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)
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)
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)
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)
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
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)
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)