Coverage for /usr/local/lib/python3.10/site-packages/opal_common-0.0.0-py3.10.egg/opal_common/cli/commands.py: 0%

70 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 11:54 +0000

1import asyncio 

2import json 

3import secrets 

4from datetime import timedelta 

5from enum import Enum 

6from typing import List, Optional, Tuple 

7from uuid import uuid4 

8 

9import typer 

10from opal_common.schemas.data import DataSourceEntry, DataUpdate 

11from opal_common.schemas.security import AccessTokenRequest, PeerType 

12 

13 

14class SecretFormat(str, Enum): 

15 hex = "hex" 

16 bytes = "bytes" 

17 urlsafe = "urlsafe" 

18 

19 

20def generate_secret( 

21 size: int = typer.Option(32, help="size in bytes of the secret"), 

22 format: SecretFormat = SecretFormat.urlsafe, 

23): 

24 if format == SecretFormat.hex: 

25 res = secrets.token_hex(size) 

26 elif format == SecretFormat.bytes: 

27 res = repr(secrets.token_bytes(size)) 

28 else: 

29 res = secrets.token_urlsafe(size) 

30 

31 typer.echo(res) 

32 

33 

34def obtain_token( 

35 master_token: str = typer.Argument( 

36 ..., 

37 help="The master token secret the OPAL-server was initialized with", 

38 envvar="OPAL_MASTER_TOKEN", 

39 ), 

40 server_url: str = typer.Option( 

41 "http://localhost:7002", help="url of the OPAL-server to obtain the token from" 

42 ), 

43 type: PeerType = PeerType("client"), 

44 ttl: Tuple[int, str] = typer.Option( 

45 (365, "days"), 

46 help="Time-To-Live / expiration for the token in `<int> <str>` e.g. `365 days`, or `1000000 milliseconds` ", 

47 ), 

48 claims: str = typer.Option( 

49 "{}", 

50 help="claims to to include in the returned signed JWT as a JSON string", 

51 callback=lambda x: json.loads(x), 

52 ), 

53 just_the_token: bool = typer.Option( 

54 True, 

55 help="Should the command return only the cryptographic token, or the full JSON object", 

56 ), 

57): 

58 """Obtain a secret JWT (JSON-Web-Token) from the server, to be used by 

59 clients or data sources for authentication Using the master token (as 

60 assigned to the server as OPAL_AUTH_MASTER_TOKEN)""" 

61 

62 from aiohttp import ClientSession 

63 

64 server_url = f"{server_url}/token" 

65 ttl_number, ttl_unit = ttl 

66 ttl = timedelta(**{ttl_unit: ttl_number}) 

67 

68 async def fetch(): 

69 async with ClientSession( 

70 headers={"Authorization": f"bearer {master_token}"}, 

71 trust_env=True, 

72 ) as session: 

73 details = AccessTokenRequest(type=type, ttl=ttl, claims=claims).json() 

74 res = await session.post( 

75 server_url, data=details, headers={"content-type": "application/json"} 

76 ) 

77 data = await res.json() 

78 if just_the_token: 

79 return data["token"] 

80 else: 

81 return data 

82 

83 res = asyncio.run(fetch()) 

84 typer.echo(res) 

85 

86 

87def publish_data_update( 

88 token: Optional[str] = typer.Argument( 

89 None, 

90 help="the JWT obtained from the server for authentication (see obtain-token command)", 

91 envvar="OPAL_CLIENT_TOKEN", 

92 ), 

93 server_url: str = typer.Option( 

94 "http://localhost:7002", 

95 help="url of the OPAL-server to send the update through", 

96 ), 

97 server_route: str = typer.Option( 

98 "/data/config", help="route in the server for update" 

99 ), 

100 reason: str = typer.Option("", help="The reason for the update"), 

101 entries: str = typer.Option( 

102 "[]", 

103 "--entries", 

104 "-e", 

105 help="Pass in the the DataUpdate entries as JSON", 

106 callback=lambda x: json.loads(x), 

107 ), 

108 src_url: str = typer.Option( 

109 None, 

110 help="[SINGLE-ENTRY-UPDATE] url of the data-source this update relates to, which the clients should approach", 

111 ), 

112 topics: List[str] = typer.Option( 

113 None, 

114 "--topic", 

115 "-t", 

116 help="[SINGLE-ENTRY-UPDATE] [List] topic (can several) for the published update (to be matched to client subscriptions)", 

117 ), 

118 data: str = typer.Option( 

119 None, 

120 help="[SINGLE-ENTRY-UPDATE] actual data to include in the update (if src_url is also supplied, it would be sent but not used)", 

121 ), 

122 src_config: str = typer.Option( 

123 "{}", 

124 help="[SINGLE-ENTRY-UPDATE] Fetching Config as JSON", 

125 callback=lambda x: json.loads(x), 

126 ), 

127 dst_path: str = typer.Option( 

128 "", 

129 help="[SINGLE-ENTRY-UPDATE] Path the client should set this value in its data-store", 

130 ), 

131 save_method: str = typer.Option( 

132 "PUT", 

133 help="[SINGLE-ENTRY-UPDATE] How the data should be saved into the give dst-path", 

134 ), 

135): 

136 """Publish a DataUpdate through an OPAL-server (indicated by --server_url). 

137 

138 [SINGLE-ENTRY-UPDATE] Send a single update DataSourceEntry via 

139 the --src-url, --src-config, --topics, --dst-path, --save-method 

140 must include --src-url to use this flow. [Multiple entries] Set 

141 DataSourceEntires as JSON (via --entries) if you include a 

142 single entry as well- it will be merged into the given JSON 

143 """ 

144 from aiohttp import ClientResponse, ClientSession 

145 

146 if not entries and not src_url: 

147 typer.secho( 

148 "You must provide either multiple entries (-e / --entries) or a single entry update (--src_url)", 

149 fg="red", 

150 ) 

151 return 

152 

153 if not isinstance(entries, list): 

154 typer.secho("Bad input for --entires was ignored", fg="red") 

155 entries = [] 

156 

157 entries: List[DataSourceEntry] 

158 

159 # single entry update (if used, we ignore the value of "entries") 

160 if src_url is not None: 

161 entries = [ 

162 DataSourceEntry( 

163 url=src_url, 

164 data=(None if data is None else json.loads(data)), 

165 topics=topics, 

166 dst_path=dst_path, 

167 save_method=save_method, 

168 config=src_config, 

169 ) 

170 ] 

171 

172 server_url = f"{server_url}{server_route}" 

173 update = DataUpdate(entries=entries, reason=reason) 

174 

175 async def publish_update(): 

176 headers = {"content-type": "application/json"} 

177 if token is not None: 

178 headers.update({"Authorization": f"bearer {token}"}) 

179 async with ClientSession(headers=headers, trust_env=True) as session: 

180 body = update.json() 

181 res = await session.post(server_url, data=body) 

182 return res 

183 

184 async def get_response_text(res: ClientResponse): 

185 return await res.text() 

186 

187 typer.echo(f"Publishing event:") 

188 typer.secho(f"{str(update)}", fg="cyan") 

189 res = asyncio.run(publish_update()) 

190 

191 if res.status == 200: 

192 typer.secho("Event Published Successfully", fg="green") 

193 else: 

194 typer.secho("Event publishing failed with status-code - {res.status}", fg="red") 

195 text = asyncio.run(get_response_text(res)) 

196 typer.echo(text) 

197 

198 

199def version(): 

200 """Print the OPAL version.""" 

201 from importlib.metadata import version 

202 

203 typer.echo(version("opal_common")) 

204 

205 

206all_commands = [obtain_token, generate_secret, publish_data_update, version]