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
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 11:54 +0000
1import asyncio
2import os
3from typing import List
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
15TOPIC_DELIMITER = "/"
16PREFIX_DELIMITER = ":"
19class DataUpdatePublisher:
20 def __init__(self, publisher: TopicPublisher) -> None:
21 self._publisher = publisher
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"]
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"]
34 Args:
35 topic (str): topic string with sub topics delimited by delimiter
37 Returns:
38 List[str]: The combinations of sub topics for the given topic
39 """
40 topic_combos = []
42 prefix = None
43 if PREFIX_DELIMITER in topic:
44 prefix, topic = topic.rsplit(":", 1)
46 sub_topics = topic.split(TOPIC_DELIMITER)
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]
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)
65 return topic_combos
67 async def publish_data_updates(self, update: DataUpdate):
68 """Notify OPAL subscribers of a new data update by topic.
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()
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 ]
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 )
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 )
112 await self._publisher.publish(
113 list(all_topic_combos), update.dict(by_alias=True)
114 )