Coverage for /usr/local/lib/python3.10/site-packages/opal_common-0.0.0-py3.10.egg/opal_common/schemas/data.py: 91%

59 statements  

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

1from typing import Any, ClassVar, Dict, List, Optional, Set, Tuple, Union 

2 

3from opal_common.fetcher.providers.http_fetch_provider import HttpFetcherConfig 

4from opal_common.logging_utils.redaction import RedactedReprMixin 

5from opal_common.schemas.store import JSONPatchAction 

6from pydantic import AnyHttpUrl, BaseModel, Field, root_validator, validator 

7 

8JsonableValue = Union[List[JSONPatchAction], List[Any], Dict[str, Any]] 

9 

10 

11DEFAULT_DATA_TOPIC = "policy_data" 

12 

13 

14class DataSourceEntry(RedactedReprMixin, BaseModel): 

15 """ 

16 Data source configuration - where client's should retrieve data from and how they should store it 

17 """ 

18 

19 # ``config`` may carry fetcher auth (e.g. Authorization headers) and 

20 # ``data`` an inline payload - mask both in repr/str so they never leak into 

21 # logs (entries are frequently interpolated into log messages). 

22 _redacted_repr_fields: ClassVar[Set[str]] = {"config", "data"} 

23 # ``url`` can embed credentials (``user:token@host`` / ``?token=``); strip 

24 # them via redact_url while keeping host/path visible for debugging. 

25 _redacted_url_fields: ClassVar[Set[str]] = {"url"} 

26 

27 @validator("data") 

28 def validate_save_method(cls, value, values): 

29 if values["save_method"] not in ["PUT", "PATCH"]: 

30 raise ValueError("'save_method' must be either PUT or PATCH") 

31 if values["save_method"] == "PATCH" and ( 31 ↛ 35line 31 didn't jump to line 35 because the condition on line 31 was never true

32 not isinstance(value, list) 

33 or not all(isinstance(elem, JSONPatchAction) for elem in value) 

34 ): 

35 raise TypeError( 

36 "'data' must be of type JSON patch request when save_method is PATCH" 

37 ) 

38 return value 

39 

40 # How to obtain the data 

41 url: str = Field(..., description="Url source to query for data") 

42 config: dict = Field( 

43 None, 

44 description="Suggested fetcher configuration (e.g. auth or method) to fetch data with", 

45 ) 

46 # How to catalog data 

47 topics: List[str] = Field( 

48 [DEFAULT_DATA_TOPIC], description="topics the data applies to" 

49 ) 

50 # How to save the data 

51 # see https://www.openpolicyagent.org/docs/latest/rest-api/#data-api path is the path nested under <OPA_SERVER>/<version>/data 

52 dst_path: str = Field("", description="OPA data api path to store the document at") 

53 save_method: str = Field( 

54 "PUT", 

55 description="Method used to write into OPA - PUT/PATCH, when using the PATCH method the data field should conform to the JSON patch schema defined in RFC 6902(https://datatracker.ietf.org/doc/html/rfc6902#section-3)", 

56 ) 

57 data: Optional[JsonableValue] = Field( 

58 None, 

59 description="Data payload to embed within the data update (instead of having " 

60 "the client fetch it from the url).", 

61 ) 

62 

63 

64class DataSourceEntryWithPollingInterval(DataSourceEntry): 

65 # Periodic Update Interval 

66 # If set, tells OPAL server how frequently to send message to clients that they need to refresh their data store from a data source 

67 # Time in Seconds 

68 periodic_update_interval: Optional[float] = Field( 

69 None, description="Polling interval to refresh data from data source" 

70 ) 

71 

72 

73class DataSourceConfig(BaseModel): 

74 """Static list of Data Source Entries returned to client. 

75 

76 Answers this question for the client: from where should i get the 

77 full picture of data i need? (as opposed to incremental data 

78 updates) 

79 """ 

80 

81 entries: List[DataSourceEntryWithPollingInterval] = Field( 

82 [], description="list of data sources and how to fetch from them" 

83 ) 

84 

85 

86class ServerDataSourceConfig(BaseModel): 

87 """As its data source configuration, the server can either hold: 

88 

89 1) A static DataSourceConfig returned to all clients regardless of 

90 identity. If all clients need the same config, this is the way to 

91 go. 

92 

93 2) A redirect url (external_source_url), to which the opal client 

94 will be redirected when requesting its DataSourceConfig. The client 

95 will issue the same request (with the same headers, including the 

96 JWT token identifying it) to the url configured. This option is good 

97 if each client must receive a different base data configuration, for 

98 example for a multi-tenant deployment. 

99 

100 By providing the server that serves external_source_url the value of 

101 OPAL_AUTH_PUBLIC_KEY, that server can validate the JWT and get it's 

102 claims, in order to apply authorization and/or other conditions 

103 before returning the data sources relevant to said client. 

104 """ 

105 

106 config: Optional[DataSourceConfig] = Field( 

107 None, description="static list of data sources and how to fetch from them" 

108 ) 

109 external_source_url: Optional[AnyHttpUrl] = Field( 

110 None, 

111 description="external url to serve data sources dynamically." 

112 + " if set, the clients will be redirected to this url when requesting to fetch data sources.", 

113 ) 

114 

115 @root_validator 

116 def check_passwords_match(cls, values): 

117 config, redirect_url = values.get("config"), values.get("external_source_url") 

118 if config is None and redirect_url is None: 118 ↛ 119line 118 didn't jump to line 119 because the condition on line 118 was never true

119 raise ValueError( 

120 "you must provide one of these fields: config, external_source_url" 

121 ) 

122 if config is not None and redirect_url is not None: 122 ↛ 123line 122 didn't jump to line 123 because the condition on line 122 was never true

123 raise ValueError( 

124 "you must provide ONLY ONE of these fields: config, external_source_url" 

125 ) 

126 return values 

127 

128 

129class CallbackEntry(BaseModel): 

130 """An entry in the callbacks register. 

131 

132 this schema is used by the callbacks api 

133 """ 

134 

135 key: Optional[str] = Field( 

136 None, description="unique id to identify this callback (optional)" 

137 ) 

138 url: str = Field(..., description="http/https url to call back on update") 

139 config: Optional[HttpFetcherConfig] = Field( 

140 None, 

141 description="optional http config for the target url (i.e: http method, headers, etc)", 

142 ) 

143 

144 

145class UpdateCallback(BaseModel): 

146 """Configuration of callbacks upon completion of a FetchEvent Allows 

147 notifying other services on the update flow. 

148 

149 Each callback is either a URL (str) or a tuple of a url and 

150 HttpFetcherConfig defining how to approach the URL 

151 """ 

152 

153 callbacks: List[Union[str, Tuple[str, HttpFetcherConfig]]] 

154 

155 

156class DataUpdate(BaseModel): 

157 """DataSources used as OPAL-server configuration Data update sent to 

158 clients.""" 

159 

160 # a UUID to identify this update (used as part of an updates complition callback) 

161 id: Optional[str] = None 

162 entries: List[DataSourceEntry] = Field( 

163 ..., description="list of related updates the OPAL client should perform" 

164 ) 

165 reason: str = Field(None, description="Reason for triggering the update") 

166 # Configuration for how to notify other services on the status of Update 

167 callback: UpdateCallback = UpdateCallback(callbacks=[]) 

168 

169 

170class DataEntryReport(BaseModel): 

171 """A report of the processing of a single DataSourceEntry.""" 

172 

173 entry: DataSourceEntry = Field(..., description="The entry that was processed") 

174 # Was the entry successfully fetched 

175 fetched: Optional[bool] = False 

176 # Was the entry successfully saved into the policy-data-store 

177 saved: Optional[bool] = False 

178 # Hash of the returned data 

179 hash: Optional[str] = None 

180 

181 

182class DataUpdateReport(BaseModel): 

183 # the UUID of the update this report is for 

184 update_id: Optional[str] = None 

185 # Each DataSourceEntry and how it was processed 

186 reports: List[DataEntryReport] 

187 # in case this is a policy update, the new hash committed the policy store. 

188 policy_hash: Optional[str] = None 

189 user_data: Dict[str, Any] = {}