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

1from functools import partial 

2from pathlib import Path 

3from typing import List, Optional 

4 

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 

21 

22 

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 ) 

52 

53 return PolicyUpdateMessageNotification(topics=topics, update=message) 

54 

55 

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 ) 

71 

72 with DiffViewer(old_commit, new_commit) as viewer: 

73 

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 

80 

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 ) 

100 

101 return PolicyUpdateMessageNotification(topics=topics, update=message) 

102 

103 

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 ) 

117 

118 if notification: 

119 async with publisher: 

120 await publisher.publish( 

121 topics=notification.topics, data=notification.update.dict() 

122 )