Coverage for pygeoapi/asyncapi.py: 27%
109 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 08:15 +0000
« 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# =================================================================
30import os
31import json
32import logging
33from pathlib import Path
34from urllib.parse import urlparse
36import click
37from jsonschema import validate as jsonschema_validate
38import yaml
40from pygeoapi import __version__, l10n
41from pygeoapi.models.openapi import OAPIFormat
42from pygeoapi.util import to_json, yaml_load, remove_url_auth
44LOGGER = logging.getLogger(__name__)
46THISDIR = os.path.dirname(os.path.realpath(__file__))
49def gen_asyncapi(cfg: dict) -> dict:
50 """
51 Generate an AsyncAPI document
53 :param cfg: `dict` of pygeoapi configuration
55 :returns: `dict` of AsyncAPI document
56 """
58 server_locales = l10n.get_locales(cfg)
59 locale_ = server_locales[0]
61 LOGGER.debug('Generating AsyncAPI document')
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
67 u = cfg['pubsub']['broker']['url']
68 up = urlparse(u)
69 protocol = up.scheme
70 url = remove_url_auth(u).replace(f'{protocol}://', '')
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']
106 if cfg['pubsub']['broker'].get('channel') is not None:
107 channel_prefix = cfg['pubsub']['broker']['channel']
108 else:
109 channel_prefix = ''
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
117 title = l10n.translate(value['title'], locale_)
118 channel_address = f'{channel_prefix}/collections/{key}'
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 }
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 }
147 a['channels'][f'notify-{key}'] = channel
148 a['operations'].update(operation)
150 return a
153def get_asyncapi(cfg, version='3.0'):
154 """
155 Stub to generate AsyncAPI Document
157 :param cfg: configuration object
158 :param version: version of AsyncAPI (default 3.0)
160 :returns: AsyncAPI definition YAML dict
161 """
163 if version == '3.0':
164 return gen_asyncapi(cfg)
165 else:
166 raise RuntimeError('AsyncAPI version not supported')
169def validate_asyncapi_document(instance_dict):
170 """
171 Validate an AsyncAPI document against the AsyncAPI schema
173 :param instance_dict: dict of AsyncAPI instance
175 :returns: `bool` of validation
176 """
178 schema_file = os.path.join(
179 THISDIR, 'resources', 'schemas', 'asyncapi', 'asyncapi-3.0.0.json')
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)
186 return True
189def generate_asyncapi_document(cfg: dict, output_format: OAPIFormat):
190 """
191 Generate an AsyncAPI document from the configuration file
193 :param cfg: `dict` of configuration
194 :param output_format: output format for AsyncAPI document
196 :returns: content of the AsyncAPI document in the output
197 format requested
198 """
200 pretty_print = cfg['server'].get('pretty_print', False)
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
209def load_asyncapi_document() -> dict:
210 """
211 Open AsyncAPI document from `PYGEOAPI_ASYNCAPI` environment variable
213 :returns: `dict` of AsyncAPI document
214 """
216 pygeoapi_asyncapi = os.environ.get('PYGEOAPI_ASYNCAPI')
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 {}
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 {}
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()
234 return asyncapi_
237@click.group()
238def asyncapi():
239 """AsyncAPI management"""
240 pass
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"""
253 if config_file is None:
254 raise click.ClickException('--config/-c required')
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)
262 if 'pubsub' not in cfg:
263 click.echo('pubsub not configured; aborting')
264 ctx.exit(1)
266 content = generate_asyncapi_document(cfg, format_)
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')
276@click.command()
277@click.pass_context
278@click.argument('asyncapi_file', type=click.File())
279def validate(ctx, asyncapi_file):
280 """Validate AsyncAPI Document"""
282 if asyncapi_file is None:
283 raise click.ClickException('--asyncapi/-o required')
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')
291asyncapi.add_command(generate)
292asyncapi.add_command(validate)