Coverage for /usr/local/lib/python3.10/site-packages/opal_server-0.0.0-py3.10.egg/opal_server/data/data_update_publisher.py: 98%

45 statements  

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

1import asyncio 

2import os 

3from typing import List 

4 

5from fastapi_utils.tasks import repeat_every 

6from opal_common.http_utils import redact_url 

7from opal_common.logger import logger 

8from opal_common.schemas.data import ( 

9 DataSourceEntryWithPollingInterval, 

10 DataUpdate, 

11 ServerDataSourceConfig, 

12) 

13from opal_common.topics.publisher import TopicPublisher 

14 

15TOPIC_DELIMITER = "/" 

16PREFIX_DELIMITER = ":" 

17 

18 

19class DataUpdatePublisher: 

20 def __init__(self, publisher: TopicPublisher) -> None: 

21 self._publisher = publisher 

22 

23 @staticmethod 

24 def get_topic_combos(topic: str) -> List[str]: 

25 """Get the The combinations of sub topics for the given topic e.g. 

26 "policy_data/users/keys" -> ["policy_data", "policy_data/users", 

27 "policy_data/users/keys"] 

28 

29 If a colon (':') is present, only split after the right-most one, 

30 and prepend the prefix before it to every topic, e.g. 

31 "data:policy_data/users/keys" -> ["data:policy_data", "data:policy_data/users", 

32 "data:policy_data/users/keys"] 

33 

34 Args: 

35 topic (str): topic string with sub topics delimited by delimiter 

36 

37 Returns: 

38 List[str]: The combinations of sub topics for the given topic 

39 """ 

40 topic_combos = [] 

41 

42 prefix = None 

43 if PREFIX_DELIMITER in topic: 

44 prefix, topic = topic.rsplit(":", 1) 

45 

46 sub_topics = topic.split(TOPIC_DELIMITER) 

47 

48 if sub_topics: 48 ↛ 65line 48 didn't jump to line 65 because the condition on line 48 was always true

49 current_topic = sub_topics[0] 

50 

51 if prefix: 

52 topic_combos.append(f"{prefix}{PREFIX_DELIMITER}{current_topic}") 

53 else: 

54 topic_combos.append(current_topic) 

55 if len(sub_topics) > 1: 

56 for sub in sub_topics[1:]: 

57 current_topic = f"{current_topic}{TOPIC_DELIMITER}{sub}" 

58 if prefix: 

59 topic_combos.append( 

60 f"{prefix}{PREFIX_DELIMITER}{current_topic}" 

61 ) 

62 else: 

63 topic_combos.append(current_topic) 

64 

65 return topic_combos 

66 

67 async def publish_data_updates(self, update: DataUpdate): 

68 """Notify OPAL subscribers of a new data update by topic. 

69 

70 Args: 

71 topics (List[str]): topics (with hierarchy) to notify subscribers of 

72 update (DataUpdate): update data-source configuration for subscribers to fetch data from 

73 """ 

74 all_topic_combos = set() 

75 

76 # a nicer format of entries to the log 

77 logged_entries = [ 

78 dict( 

79 url=redact_url(entry.url), 

80 method=entry.save_method, 

81 path=entry.dst_path or "/", 

82 inline_data=(entry.data is not None), 

83 topics=entry.topics, 

84 ) 

85 for entry in update.entries 

86 ] 

87 

88 # Expand the topics for each event to include sub topic combos (e.g. publish 'a/b/c' as 'a' , 'a/b', and 'a/b/c') 

89 for entry in update.entries: 

90 topic_combos = [] 

91 if entry.topics: 

92 for topic in entry.topics: 

93 topic_combos.extend(DataUpdatePublisher.get_topic_combos(topic)) 

94 entry.topics = topic_combos # Update entry with the exhaustive list, so client won't have to expand it again 

95 all_topic_combos.update(topic_combos) 

96 else: 

97 logger.warning( 

98 "[{pid}] No topics were provided for the entry with url: {url}", 

99 pid=os.getpid(), 

100 url=redact_url(entry.url), 

101 ) 

102 

103 # publish all topics with all their sub combinations 

104 logger.info( 

105 "[{pid}] Publishing data update to topics: {topics}, reason: {reason}, entries: {entries}", 

106 pid=os.getpid(), 

107 topics=all_topic_combos, 

108 reason=update.reason, 

109 entries=logged_entries, 

110 ) 

111 

112 await self._publisher.publish( 

113 list(all_topic_combos), update.dict(by_alias=True) 

114 )