Coverage for pygeoapi/asyncapi.py: 27%

109 statements  

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

1# ================================================================= 

2# 

3# Authors: Tom Kralidis <tomkralidis@gmail.com> 

4# 

5# Copyright (c) 2026 Tom Kralidis 

6# 

7# Permission is hereby granted, free of charge, to any person 

8# obtaining a copy of this software and associated documentation 

9# files (the "Software"), to deal in the Software without 

10# restriction, including without limitation the rights to use, 

11# copy, modify, merge, publish, distribute, sublicense, and/or sell 

12# copies of the Software, and to permit persons to whom the 

13# Software is furnished to do so, subject to the following 

14# conditions: 

15# 

16# The above copyright notice and this permission notice shall be 

17# included in all copies or substantial portions of the Software. 

18# 

19# THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, 

20# EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES 

21# OF MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND 

22# NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR COPYRIGHT 

23# HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, 

24# WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING 

25# FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR 

26# OTHER DEALINGS IN THE SOFTWARE. 

27# 

28# ================================================================= 

29 

30import os 

31import json 

32import logging 

33from pathlib import Path 

34from urllib.parse import urlparse 

35 

36import click 

37from jsonschema import validate as jsonschema_validate 

38import yaml 

39 

40from pygeoapi import __version__, l10n 

41from pygeoapi.models.openapi import OAPIFormat 

42from pygeoapi.util import to_json, yaml_load, remove_url_auth 

43 

44LOGGER = logging.getLogger(__name__) 

45 

46THISDIR = os.path.dirname(os.path.realpath(__file__)) 

47 

48 

49def gen_asyncapi(cfg: dict) -> dict: 

50 """ 

51 Generate an AsyncAPI document 

52 

53 :param cfg: `dict` of pygeoapi configuration 

54 

55 :returns: `dict` of AsyncAPI document 

56 """ 

57 

58 server_locales = l10n.get_locales(cfg) 

59 locale_ = server_locales[0] 

60 

61 LOGGER.debug('Generating AsyncAPI document') 

62 

63 title = l10n.translate(cfg['metadata']['identification']['title'], locale_) # noqa 

64 description = l10n.translate(cfg['metadata']['identification']['description'], locale_) # noqa 

65 tags = l10n.translate(cfg['metadata']['identification']['keywords'], locale_) # noqa 

66 

67 u = cfg['pubsub']['broker']['url'] 

68 up = urlparse(u) 

69 protocol = up.scheme 

70 url = remove_url_auth(u).replace(f'{protocol}://', '') 

71 

72 a = { 

73 'asyncapi': '3.0.0', 

74 'id': cfg['server']['url'], 

75 'defaultContentType': 'application/json', 

76 'info': { 

77 'version': __version__, 

78 'title': title, 

79 'description': description, 

80 'license': { 

81 'name': cfg['metadata']['license']['name'], 

82 'url': cfg['metadata']['license']['url'] 

83 }, 

84 'contact': { 

85 'name': cfg['metadata']['contact']['name'], 

86 'email': cfg['metadata']['contact']['email'] 

87 }, 

88 'tags': [{'name': tag} for tag in tags], 

89 'externalDocs': { 

90 'url': cfg['metadata']['identification']['url'] 

91 }, 

92 }, 

93 'servers': { 

94 'default': { 

95 'host': url, 

96 'protocol': protocol, 

97 'description': description 

98 } 

99 }, 

100 'channels': {}, 

101 'operations': {} 

102 } 

103 if cfg['metadata']['contact']['url'].startswith('http'): 

104 a['info']['contact']['url'] = cfg['metadata']['contact']['url'] 

105 

106 if cfg['pubsub']['broker'].get('channel') is not None: 

107 channel_prefix = cfg['pubsub']['broker']['channel'] 

108 else: 

109 channel_prefix = '' 

110 

111 LOGGER.debug('Generating channels foreach collection') 

112 for key, value in cfg['resources'].items(): 

113 if value['type'] not in ['collection']: 

114 LOGGER.debug('Skipping') 

115 continue 

116 

117 title = l10n.translate(value['title'], locale_) 

118 channel_address = f'{channel_prefix}/collections/{key}' 

119 

120 channel = { 

121 'description': title, 

122 'address': channel_address, 

123 'messages': { 

124 'DefaultMessage': { 

125 'payload': { 

126 '$ref': 'https://raw.githubusercontent.com/wmo-im/wis2-monitoring-events/refs/heads/main/schemas/cloudevents-v1.0.2.yaml' # noqa 

127 } 

128 } 

129 } 

130 } 

131 

132 operation = { 

133 f'publish-{key}': { 

134 'action': 'send', 

135 'channel': { 

136 '$ref': f'#/channels/notify-{key}' 

137 } 

138 }, 

139 f'consume-{key}': { 

140 'action': 'receive', 

141 'channel': { 

142 '$ref': f'#/channels/notify-{key}' 

143 } 

144 } 

145 } 

146 

147 a['channels'][f'notify-{key}'] = channel 

148 a['operations'].update(operation) 

149 

150 return a 

151 

152 

153def get_asyncapi(cfg, version='3.0'): 

154 """ 

155 Stub to generate AsyncAPI Document 

156 

157 :param cfg: configuration object 

158 :param version: version of AsyncAPI (default 3.0) 

159 

160 :returns: AsyncAPI definition YAML dict 

161 """ 

162 

163 if version == '3.0': 

164 return gen_asyncapi(cfg) 

165 else: 

166 raise RuntimeError('AsyncAPI version not supported') 

167 

168 

169def validate_asyncapi_document(instance_dict): 

170 """ 

171 Validate an AsyncAPI document against the AsyncAPI schema 

172 

173 :param instance_dict: dict of AsyncAPI instance 

174 

175 :returns: `bool` of validation 

176 """ 

177 

178 schema_file = os.path.join( 

179 THISDIR, 'resources', 'schemas', 'asyncapi', 'asyncapi-3.0.0.json') 

180 

181 LOGGER.debug(f'Validating against {schema_file}') 

182 with open(schema_file) as fh2: 

183 schema_dict = json.load(fh2) 

184 jsonschema_validate(instance_dict, schema_dict) 

185 

186 return True 

187 

188 

189def generate_asyncapi_document(cfg: dict, output_format: OAPIFormat): 

190 """ 

191 Generate an AsyncAPI document from the configuration file 

192 

193 :param cfg: `dict` of configuration 

194 :param output_format: output format for AsyncAPI document 

195 

196 :returns: content of the AsyncAPI document in the output 

197 format requested 

198 """ 

199 

200 pretty_print = cfg['server'].get('pretty_print', False) 

201 

202 if output_format == 'yaml': 

203 content = yaml.safe_dump(get_asyncapi(cfg), default_flow_style=False) 

204 else: 

205 content = to_json(get_asyncapi(cfg), pretty=pretty_print) 

206 return content 

207 

208 

209def load_asyncapi_document() -> dict: 

210 """ 

211 Open AsyncAPI document from `PYGEOAPI_ASYNCAPI` environment variable 

212 

213 :returns: `dict` of AsyncAPI document 

214 """ 

215 

216 pygeoapi_asyncapi = os.environ.get('PYGEOAPI_ASYNCAPI') 

217 

218 if pygeoapi_asyncapi is None: 218 ↛ 222line 218 didn't jump to line 222 because the condition on line 218 was always true

219 LOGGER.debug('PYGEOAPI_ASYNCAPI environment not set') 

220 return {} 

221 

222 if not os.path.exists(pygeoapi_asyncapi): 

223 msg = (f'AsyncAPI document {pygeoapi_asyncapi} does not exist. ' 

224 'Please generate before starting pygeoapi') 

225 LOGGER.warning(msg) 

226 return {} 

227 

228 with open(pygeoapi_asyncapi, encoding='utf8') as ff: 

229 if pygeoapi_asyncapi.endswith(('.yaml', '.yml')): 

230 asyncapi_ = yaml_load(ff) 

231 else: # JSON string, do not transform 

232 asyncapi_ = ff.read() 

233 

234 return asyncapi_ 

235 

236 

237@click.group() 

238def asyncapi(): 

239 """AsyncAPI management""" 

240 pass 

241 

242 

243@click.command() 

244@click.pass_context 

245@click.argument('config_file', type=click.File(encoding='utf-8')) 

246@click.option('--format', '-f', 'format_', type=click.Choice(['json', 'yaml']), 

247 default='yaml', help='output format (json|yaml)') 

248@click.option('--output-file', '-of', type=click.File('w', encoding='utf-8'), 

249 help='Name of output file') 

250def generate(ctx, config_file, output_file, format_='yaml'): 

251 """Generate AsyncAPI Document""" 

252 

253 if config_file is None: 

254 raise click.ClickException('--config/-c required') 

255 

256 if isinstance(config_file, Path): 

257 with config_file.open(mode='r') as cf: 

258 cfg = yaml_load(cf) 

259 else: 

260 cfg = yaml_load(config_file) 

261 

262 if 'pubsub' not in cfg: 

263 click.echo('pubsub not configured; aborting') 

264 ctx.exit(1) 

265 

266 content = generate_asyncapi_document(cfg, format_) 

267 

268 if output_file is None: 

269 click.echo(content) 

270 else: 

271 click.echo(f'Generating {output_file.name}') 

272 output_file.write(content) 

273 click.echo('Done') 

274 

275 

276@click.command() 

277@click.pass_context 

278@click.argument('asyncapi_file', type=click.File()) 

279def validate(ctx, asyncapi_file): 

280 """Validate AsyncAPI Document""" 

281 

282 if asyncapi_file is None: 

283 raise click.ClickException('--asyncapi/-o required') 

284 

285 click.echo(f'Validating {asyncapi_file.name}') 

286 instance = yaml_load(asyncapi_file) 

287 validate_asyncapi_document(instance) 

288 click.echo('Valid AsyncAPI document') 

289 

290 

291asyncapi.add_command(generate) 

292asyncapi.add_command(validate)