Coverage for oc_meta / core / editor.py: 67%
162 statements
« prev ^ index » next coverage.py v7.13.4, created at 2026-07-25 10:39 +0000
« prev ^ index » next coverage.py v7.13.4, created at 2026-07-25 10:39 +0000
1#!/usr/bin/python
3# SPDX-FileCopyrightText: 2026 Arcangelo Massari <arcangelo.massari@unibo.it>
4#
5# SPDX-License-Identifier: ISC
7from __future__ import annotations
9import os
11import validators
12import yaml
13from oc_ocdm import Storer
14from oc_ocdm.counter_handler.counter_handler import CounterHandler
15from oc_ocdm.counter_handler.filesystem_counter_handler import FilesystemCounterHandler
16from oc_ocdm.graph import GraphSet
17from oc_ocdm.graph.graph_entity import GraphEntity
18from triplelite import SubgraphView
19from oc_ocdm.prov import ProvSet
20from oc_ocdm.reader import Reader
21from oc_ocdm.support import get_prefix
22from oc_ocdm.support.support import build_graph_from_results
24from oc_meta.lib.file_manager import find_rdf_file
25from oc_meta.lib.sparql import execute_sparql
28class EntityCache:
29 def __init__(self):
30 self.cache: set[str] = set()
32 def add(self, entity: str) -> None:
33 self.cache.add(entity)
35 def is_cached(self, entity: str) -> bool:
36 return entity in self.cache
38 def clear(self) -> None:
39 self.cache.clear()
42class MetaEditor:
43 property_to_remove_method = {
44 "http://purl.org/spar/datacite/hasIdentifier": "remove_identifier",
45 "http://purl.org/spar/pro/isHeldBy": "remove_is_held_by",
46 "http://purl.org/vocab/frbr/core#embodiment": "remove_format",
47 "http://purl.org/spar/pro/isDocumentContextFor": "remove_is_held_by",
48 "https://w3id.org/oc/ontology/hasNext": "remove_next",
49 }
51 def __init__(
52 self,
53 meta_config: str,
54 resp_agent: str,
55 save_queries: bool = False,
56 counter_handler: CounterHandler | None = None,
57 ):
58 with open(meta_config, encoding="utf-8") as file:
59 settings = yaml.full_load(file)
60 self.endpoint = settings["triplestore_url"]
61 self.provenance_endpoint = settings["provenance_triplestore_url"]
62 output_dir = settings.get("base_output_dir")
63 self.data_hotfix_dir = os.path.join(output_dir, "to_be_uploaded_hotfix")
64 self.prov_hotfix_dir = os.path.join(output_dir, "to_be_uploaded_hotfix")
65 self.base_dir = os.path.join(output_dir, "rdf") + os.sep
66 self.base_iri = settings["base_iri"]
67 self.resp_agent = resp_agent
68 self.dir_split = settings["dir_split_number"]
69 self.n_file_item = settings["items_per_file"]
70 self.zip_output_rdf = settings["zip_output_rdf"]
71 self.rdf_files_only = settings.get("rdf_files_only", False)
72 self.reader = Reader()
73 self.save_queries = save_queries
74 self.update_queries = []
76 supplier_prefix = (
77 settings["supplier_prefix"] if "supplier_prefix" in settings else "060"
78 )
79 if not supplier_prefix.endswith("0"):
80 supplier_prefix = f"{supplier_prefix}0"
81 self.supplier_prefix = supplier_prefix
83 if counter_handler is not None:
84 self.counter_handler = counter_handler
85 else:
86 info_dir = os.path.join(output_dir, "info_dir", supplier_prefix) + os.sep
87 self.counter_handler = FilesystemCounterHandler(
88 info_dir=info_dir, supplier_prefix=supplier_prefix
89 )
91 self.entity_cache = EntityCache()
93 def update_property(self, res: str, property: str, new_value: str) -> None:
94 supplier_prefix = get_prefix(res)
95 g_set = GraphSet(
96 self.base_iri,
97 supplier_prefix=supplier_prefix,
98 custom_counter_handler=self.counter_handler,
99 )
100 self.reader.import_entity_from_triplestore(
101 g_set, self.endpoint, res, self.resp_agent, enable_validation=False
102 )
103 if validators.url(new_value):
104 self.reader.import_entity_from_triplestore(
105 g_set,
106 self.endpoint,
107 new_value,
108 self.resp_agent,
109 enable_validation=False,
110 )
111 getattr(g_set.get_entity(res), property)(g_set.get_entity(new_value))
112 else:
113 getattr(g_set.get_entity(res), property)(new_value)
114 self.save(g_set, supplier_prefix)
116 def delete(
117 self, res: str, property: str | None = None, object: str | None = None
118 ) -> None:
119 res_str = str(res)
120 supplier_prefix = get_prefix(res_str)
121 g_set = GraphSet(
122 self.base_iri,
123 supplier_prefix=supplier_prefix,
124 custom_counter_handler=self.counter_handler,
125 )
126 try:
127 self.reader.import_entity_from_triplestore(
128 g_set, self.endpoint, res_str, self.resp_agent, enable_validation=False
129 )
130 except ValueError:
131 inferred_type = self.infer_type_from_uri(res_str)
132 if inferred_type:
133 query: str = (
134 f"SELECT ?s ?p ?o WHERE {{BIND (<{res_str}> AS ?s). ?s ?p ?o.}}"
135 )
136 result = execute_sparql(
137 self.endpoint, query, max_retries=3, backoff_factor=0.3
138 )["results"]["bindings"]
139 graph = build_graph_from_results(result)
140 preexisting_graph = graph.subgraph(res_str)
141 self.add_entity_with_type(
142 g_set, res_str, inferred_type, preexisting_graph
143 )
144 else:
145 return
146 if not g_set.get_entity(res_str):
147 return
148 if property:
149 remove_method = (
150 self.property_to_remove_method[property]
151 if property in self.property_to_remove_method
152 else (
153 property.replace("has", "remove")
154 if property.startswith("has")
155 else f"remove_{property}"
156 )
157 )
158 if object:
159 if validators.url(object):
160 self.reader.import_entity_from_triplestore(
161 g_set,
162 self.endpoint,
163 object,
164 self.resp_agent,
165 enable_validation=False,
166 )
167 getattr(g_set.get_entity(res_str), remove_method)(
168 g_set.get_entity(object)
169 )
170 else:
171 getattr(g_set.get_entity(res_str), remove_method)(object)
172 else:
173 getattr(g_set.get_entity(res_str), remove_method)()
174 else:
175 query = f"SELECT ?s WHERE {{?s ?p <{res_str}>.}}"
176 result = execute_sparql(
177 self.endpoint, query, max_retries=3, backoff_factor=0.3
178 )
179 for entity in result["results"]["bindings"]:
180 self.reader.import_entity_from_triplestore(
181 g_set,
182 self.endpoint,
183 entity["s"]["value"],
184 self.resp_agent,
185 enable_validation=False,
186 )
187 entity_to_purge = g_set.get_entity(res_str)
188 if not entity_to_purge:
189 return
190 entity_to_purge.mark_as_to_be_deleted()
191 self.save(g_set, supplier_prefix)
193 def sync_rdf_with_triplestore(
194 self, res: str, source_uri: str | None = None
195 ) -> bool:
196 supplier_prefix = get_prefix(res)
197 g_set = GraphSet(
198 self.base_iri,
199 supplier_prefix=supplier_prefix,
200 custom_counter_handler=self.counter_handler,
201 )
202 try:
203 self.reader.import_entity_from_triplestore(
204 g_set, self.endpoint, res, self.resp_agent, enable_validation=False
205 )
206 self.save(g_set, supplier_prefix)
207 return True
208 except ValueError:
209 if not source_uri:
210 return False
211 try:
212 self.reader.import_entity_from_triplestore(
213 g_set,
214 self.endpoint,
215 source_uri,
216 self.resp_agent,
217 enable_validation=False,
218 )
219 return False
220 except ValueError:
221 res_filepath = find_rdf_file(
222 source_uri,
223 self.base_dir,
224 self.dir_split,
225 self.n_file_item,
226 self.zip_output_rdf,
227 )
228 if not res_filepath:
229 return False
230 imported_graph = self.reader.load(res_filepath)
231 if not imported_graph:
232 return False
233 self.reader.import_entities_from_graph(
234 g_set, imported_graph, self.resp_agent
235 )
236 res_entity = g_set.get_entity(source_uri)
237 if res_entity:
238 for entity_res, entity in g_set.res_to_entity.items():
239 triples_list = list(entity.g.triples((source_uri, None, None)))
240 for triple in triples_list:
241 entity.g.remove(triple)
242 self.save(g_set, supplier_prefix)
243 return False
245 def save(self, g_set: GraphSet, supplier_prefix: str = "") -> None:
246 provset = ProvSet(
247 g_set,
248 self.base_iri,
249 wanted_label=False,
250 supplier_prefix=supplier_prefix,
251 custom_counter_handler=self.counter_handler,
252 )
253 provset.generate_provenance()
254 graph_storer = Storer(
255 g_set,
256 dir_split=self.dir_split,
257 n_file_item=self.n_file_item,
258 zip_output=self.zip_output_rdf,
259 )
260 prov_storer = Storer(
261 provset,
262 dir_split=self.dir_split,
263 n_file_item=self.n_file_item,
264 zip_output=self.zip_output_rdf,
265 )
267 graph_storer.store_all(self.base_dir, self.base_iri)
268 prov_storer.store_all(self.base_dir, self.base_iri)
270 if not self.rdf_files_only:
271 graph_storer.upload_all(
272 self.endpoint,
273 base_dir=self.data_hotfix_dir,
274 save_queries=self.save_queries,
275 )
276 prov_storer.upload_all(
277 self.provenance_endpoint,
278 base_dir=self.prov_hotfix_dir,
279 save_queries=self.save_queries,
280 )
281 g_set.commit_changes()
282 if isinstance(self.counter_handler, FilesystemCounterHandler):
283 self.counter_handler.flush()
285 def infer_type_from_uri(self, uri: str) -> str | None:
286 if os.path.join(self.base_iri, "br") in uri:
287 return GraphEntity.iri_expression
288 elif os.path.join(self.base_iri, "ar") in uri:
289 return GraphEntity.iri_role_in_time
290 elif os.path.join(self.base_iri, "ra") in uri:
291 return GraphEntity.iri_agent
292 elif os.path.join(self.base_iri, "re") in uri:
293 return GraphEntity.iri_manifestation
294 elif os.path.join(self.base_iri, "id") in uri:
295 return GraphEntity.iri_identifier
296 return None
298 def add_entity_with_type(
299 self,
300 g_set: GraphSet,
301 res: str,
302 entity_type: str,
303 preexisting_graph: SubgraphView | None,
304 ):
305 if entity_type == GraphEntity.iri_expression:
306 g_set.add_br(
307 resp_agent=self.resp_agent,
308 res=res,
309 preexisting_graph=preexisting_graph,
310 )
311 elif entity_type == GraphEntity.iri_role_in_time:
312 g_set.add_ar(
313 resp_agent=self.resp_agent,
314 res=res,
315 preexisting_graph=preexisting_graph,
316 )
317 elif entity_type == GraphEntity.iri_agent:
318 g_set.add_ra(
319 resp_agent=self.resp_agent,
320 res=res,
321 preexisting_graph=preexisting_graph,
322 )
323 elif entity_type == GraphEntity.iri_manifestation:
324 g_set.add_re(
325 resp_agent=self.resp_agent,
326 res=res,
327 preexisting_graph=preexisting_graph,
328 )
329 elif entity_type == GraphEntity.iri_identifier:
330 g_set.add_id(
331 resp_agent=self.resp_agent,
332 res=res,
333 preexisting_graph=preexisting_graph,
334 )