Coverage for pygeoapi/api/pubsub.py: 21%
43 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# =================================================================
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# =================================================================
30from datetime import datetime, UTC
31import json
32import logging
33import uuid
34from typing import Union
36LOGGER = logging.getLogger(__name__)
38CONFORMANCE_CLASSES = [
39 'https://www.opengis.net/spec/ogcapi-pubsub-1/1.0/conf/message-payload-cloudevents-json', # noqa
40 'https://www.opengis.net/spec/ogcapi-pubsub-1/1.0/conf/discovery'
41]
44def publish_message(pubsub_client, url: str, action: str,
45 resource: str = None, item: str = None,
46 data: dict = None) -> bool:
47 """
48 Publish broker message
50 :param pubsub_client: `pygeoapi.pubsub.BasePubSubClient` instance
51 :param url: `str` of server base URL
52 :param action: `str` of action trigger name (create, update, delete)
53 :param resource: `str` of resource identifier
54 :param item: `str` of item identifier
55 :param data: `dict` of data payload
57 :returns: `bool` of whether message publishing was successful
58 """
60 if action in ['create', 'update']:
61 channel = f'collections/{resource}'
62 data_ = data
63 media_type = 'application/geo+json'
64 type_ = f'org.ogc.api.collection.item.{action}'
65 elif action == 'delete':
66 channel = f'collections/{resource}'
67 data_ = item
68 media_type = 'text/plain'
69 type_ = f'org.ogc.api.collection.item.{action}'
70 elif action == 'process':
71 channel = f'processes/{resource}'
72 media_type = 'application/json'
73 data_ = data
74 type_ = 'org.ogc.api.job.result'
76 if pubsub_client.channel is not None:
77 channel = f'{pubsub_client.channel}/{channel}'
79 message = generate_ogc_cloudevent(type_, media_type, url,
80 channel, data_)
81 LOGGER.debug(f'Message: {message}')
83 try:
84 pubsub_client.connect()
85 pubsub_client.pub(channel, json.dumps(message))
86 except Exception as err:
87 raise RuntimeError(err)
90def generate_ogc_cloudevent(type_: str, media_type: str, source: str,
91 subject: str, data: Union[dict, str]) -> dict:
92 """
93 Generate CloudEvent
95 :param type_: `str` of CloudEvents type
96 :param source: `str` of source
97 :param subject: `str` of subject
98 :param media_type: `str` of media type
99 :param data: `str` or `dict` of data
101 :returns: `dict` of OGC CloudEvent payload
102 """
104 try:
105 data2 = json.loads(data)
106 except Exception:
107 if isinstance(data, bytes):
108 data2 = data.decode('utf-8')
109 else:
110 data2 = data
112 message = {
113 'specversion': '1.0',
114 'type': type_,
115 'source': source,
116 'subject': subject,
117 'id': str(uuid.uuid4()),
118 'time': datetime.now(UTC).strftime('%Y-%m-%dT%H:%M:%SZ'),
119 'datacontenttype': media_type,
120 # 'dataschema': 'TODO',
121 'data': data2
122 }
124 return message
127def get_oas_30(cfg, locale_):
128 return [], {}