Coverage for /usr/local/lib/python3.10/site-packages/opal_server-0.0.0-py3.10.egg/opal_server/policy/watcher/callbacks.py: 40%
48 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
1from functools import partial
2from pathlib import Path
3from typing import List, Optional
5from git.objects import Commit
6from opal_common.git_utils.commit_viewer import (
7 CommitViewer,
8 FileFilter,
9 find_ignore_match,
10 has_extension,
11)
12from opal_common.git_utils.diff_viewer import DiffViewer
13from opal_common.logger import logger
14from opal_common.paths import PathUtils
15from opal_common.schemas.policy import (
16 PolicyUpdateMessage,
17 PolicyUpdateMessageNotification,
18)
19from opal_common.topics.publisher import TopicPublisher
20from opal_common.topics.utils import policy_topics
23async def create_update_all_directories_in_repo(
24 old_commit: Commit,
25 new_commit: Commit,
26 file_extensions: Optional[List[str]] = None,
27 bundle_ignore: Optional[List[str]] = None,
28 predicate: Optional[FileFilter] = None,
29) -> PolicyUpdateMessageNotification:
30 """Publishes policy topics matching all relevant directories in tracked
31 repo, prompting the client to ask for *all* contents of these directories
32 (and not just diffs)."""
33 with CommitViewer(new_commit) as viewer:
34 if predicate is None: 34 ↛ 35line 34 didn't jump to line 35 because the condition on line 34 was never true
35 _has_extension = partial(has_extension, extensions=file_extensions)
36 _find_ignore_match = partial(find_ignore_match, bundle_ignore=bundle_ignore)
37 filter = lambda f: _has_extension(f) and _find_ignore_match(f.path) == None
38 else:
39 filter = predicate
40 all_paths = [p.path for p in list(viewer.files(filter))]
41 directories = PathUtils.intermediate_directories(all_paths)
42 logger.info(
43 "Publishing policy update, directories: {directories}",
44 directories=[str(d) for d in directories],
45 )
46 topics = policy_topics(directories)
47 message = PolicyUpdateMessage(
48 old_policy_hash=old_commit.hexsha,
49 new_policy_hash=new_commit.hexsha,
50 changed_directories=[str(path) for path in directories],
51 )
53 return PolicyUpdateMessageNotification(topics=topics, update=message)
56async def create_policy_update(
57 old_commit: Commit,
58 new_commit: Commit,
59 file_extensions: Optional[List[str]] = None,
60 bundle_ignore: Optional[List[str]] = None,
61 predicate: Optional[FileFilter] = None,
62) -> Optional[PolicyUpdateMessageNotification]:
63 if new_commit == old_commit:
64 return await create_update_all_directories_in_repo(
65 old_commit,
66 new_commit,
67 file_extensions=file_extensions,
68 bundle_ignore=bundle_ignore,
69 predicate=predicate,
70 )
72 with DiffViewer(old_commit, new_commit) as viewer:
74 def is_path_affected(path: Path) -> bool:
75 if not file_extensions:
76 return True
77 if not path.suffix in file_extensions:
78 return False
79 return find_ignore_match(path, bundle_ignore) is None
81 all_paths = list(viewer.affected_paths(is_path_affected))
82 if not all_paths:
83 logger.warning(
84 f"new commits detected but no tracked files were affected: '{old_commit.hexsha}' -> '{new_commit.hexsha}'",
85 old_commit=old_commit,
86 new_commit=new_commit,
87 )
88 return None
89 directories = PathUtils.intermediate_directories(all_paths)
90 logger.debug(
91 "Generating policy update notification, directories: {directories}",
92 directories=[str(d) for d in directories],
93 )
94 topics = policy_topics(directories)
95 message = PolicyUpdateMessage(
96 old_policy_hash=old_commit.hexsha,
97 new_policy_hash=new_commit.hexsha,
98 changed_directories=[str(path) for path in directories],
99 )
101 return PolicyUpdateMessageNotification(topics=topics, update=message)
104async def publish_changed_directories(
105 old_commit: Commit,
106 new_commit: Commit,
107 publisher: TopicPublisher,
108 file_extensions: Optional[List[str]] = None,
109 bundle_ignore: Optional[List[str]] = None,
110):
111 """Publishes policy topics matching all relevant directories in tracked
112 repo, prompting the client to ask for *all* contents of these directories
113 (and not just diffs)."""
114 notification = await create_policy_update(
115 old_commit, new_commit, file_extensions, bundle_ignore
116 )
118 if notification:
119 async with publisher:
120 await publisher.publish(
121 topics=notification.topics, data=notification.update.dict()
122 )