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

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 

30from datetime import datetime, UTC 

31import json 

32import logging 

33import uuid 

34from typing import Union 

35 

36LOGGER = logging.getLogger(__name__) 

37 

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] 

42 

43 

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 

49 

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 

56 

57 :returns: `bool` of whether message publishing was successful 

58 """ 

59 

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' 

75 

76 if pubsub_client.channel is not None: 

77 channel = f'{pubsub_client.channel}/{channel}' 

78 

79 message = generate_ogc_cloudevent(type_, media_type, url, 

80 channel, data_) 

81 LOGGER.debug(f'Message: {message}') 

82 

83 try: 

84 pubsub_client.connect() 

85 pubsub_client.pub(channel, json.dumps(message)) 

86 except Exception as err: 

87 raise RuntimeError(err) 

88 

89 

90def generate_ogc_cloudevent(type_: str, media_type: str, source: str, 

91 subject: str, data: Union[dict, str]) -> dict: 

92 """ 

93 Generate CloudEvent 

94 

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 

100 

101 :returns: `dict` of OGC CloudEvent payload 

102 """ 

103 

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 

111 

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 } 

123 

124 return message 

125 

126 

127def get_oas_30(cfg, locale_): 

128 return [], {}