Coverage for src/time_agnostic_library/agnostic_entity.py: 99%
594 statements
« prev ^ index » next coverage.py v7.15.4, created at 2026-09-03 21:17 +0000
« prev ^ index » next coverage.py v7.15.4, created at 2026-09-03 21:17 +0000
1# SPDX-FileCopyrightText: 2021-2026 Arcangelo Massari <arcangelo.massari@unibo.it>
2#
3# SPDX-License-Identifier: ISC
5import re
6from collections.abc import Callable, Iterator
8from time_agnostic_library.prov_entity import ProvEntity
9from time_agnostic_library.sparql import Sparql, _n3_value
10from time_agnostic_library.support import _cached_parse as _parse_datetime
11from time_agnostic_library.support import convert_to_datetime
13# The alternatives before the keyword consume the constructs a keyword can hide
14# inside, so that only one standing in the query itself is reported.
15_UPDATE_TOKEN_RE = re.compile(
16 r"<[^>]*>" # IRI
17 r'|"(?:[^"\\]|\\.)*"' # literal between double quotes
18 r"|'(?:[^'\\]|\\.)*'" # literal between single quotes
19 r"|#[^\r\n]*" # comment, up to the end of the line
20 r"|(?P<operation>DELETE|INSERT)\s+DATA", # the keyword being looked for
21 re.IGNORECASE,
22)
23_GRAPH_BLOCK_RE = re.compile(r"GRAPH\s*<([^>]+)>\s*\{", re.IGNORECASE)
25_RDF_TERM_RE = re.compile(
26 r"<(?P<iri>[^>]+)>"
27 r'|"(?P<typed>(?:[^"\\]|\\.)*)"\^\^<(?P<datatype>[^>]+)>'
28 r'|"(?P<tagged>(?:[^"\\]|\\.)*)"@(?P<language>[a-zA-Z][\w-]*)'
29 r'|"(?P<quoted>(?:[^"\\]|\\.)*)"'
30 r"|'(?P<single_quoted>(?:[^'\\]|\\.)*)'"
31 r"|(?P<blank_node>_:\S+)",
32 re.DOTALL,
33)
35_ESCAPE_CHAR_RE = re.compile(r"\\(.)")
36_ESCAPE_CHAR_MAP = {"n": "\n", "r": "\r", "t": "\t"}
38_RDF_TYPE = "http://www.w3.org/1999/02/22-rdf-syntax-ns#type"
40_TRIPLE_LEN = 3
43def _unescape_literal(s: str) -> str:
44 if "\\" not in s:
45 return s
46 return _ESCAPE_CHAR_RE.sub(
47 lambda m: _ESCAPE_CHAR_MAP.get(m.group(1), m.group(1)), s
48 )
51def _normalize_literal(raw: str) -> str:
52 unescaped = _unescape_literal(raw)
53 return (
54 unescaped.replace("\\", "\\\\")
55 .replace('"', '\\"')
56 .replace("\n", "\\n")
57 .replace("\r", "\\r")
58 )
61def _regex_match_to_n3(match: re.Match) -> str:
62 iri = match.group("iri")
63 if iri is not None:
64 return f"<{iri}>"
66 typed = match.group("typed")
67 if typed is not None:
68 return f'"{_normalize_literal(typed)}"^^<{match.group("datatype")}>'
70 tagged = match.group("tagged")
71 if tagged is not None:
72 return f'"{_normalize_literal(tagged)}"@{match.group("language")}'
74 quoted = match.group("quoted")
75 if quoted is not None:
76 return f'"{_normalize_literal(quoted)}"'
78 single_quoted = match.group("single_quoted")
79 if single_quoted is not None:
80 return f'"{_normalize_literal(single_quoted)}"'
82 return match.group("blank_node")
85_STRUCTURE_TOKEN_RE = re.compile(
86 r"<[^>]*>" # IRI
87 r'|"(?:[^"\\]|\\.)*"' # literal between double quotes
88 r"|'(?:[^'\\]|\\.)*'" # literal between single quotes
89 r"|#[^\r\n]*" # comment, up to the end of the line
90 r"|(?P<brace>[{}])"
91)
94def _find_matching_close_brace(text: str, start: int) -> int:
95 depth = 1
96 for token_match in _STRUCTURE_TOKEN_RE.finditer(text, start):
97 brace = token_match.group("brace")
98 if brace is None:
99 continue
100 if brace == "{":
101 depth += 1
102 continue
103 depth -= 1
104 if depth == 0:
105 return token_match.start()
106 return len(text)
109def _parse_graph_blocks(
110 text: str, start: int, end: int
111) -> list[tuple[str, str, str, str]]:
112 quads: list[tuple[str, str, str, str]] = []
113 pos = start
114 while pos < end:
115 graph_match = _GRAPH_BLOCK_RE.search(text, pos, end)
116 if graph_match is None:
117 break
118 graph_n3 = f"<{graph_match.group(1)}>"
119 triples_start = graph_match.end()
120 triples_end = min(_find_matching_close_brace(text, triples_start), end)
122 terms: list[str] = []
123 for m in _RDF_TERM_RE.finditer(text, triples_start, triples_end):
124 terms.append(_regex_match_to_n3(m))
125 if len(terms) == _TRIPLE_LEN:
126 quads.append((terms[0], terms[1], terms[2], graph_n3))
127 terms.clear()
129 pos = triples_end + 1
130 return quads
133def _find_next_operation(text: str, start: int) -> re.Match | None:
134 for token_match in _UPDATE_TOKEN_RE.finditer(text, start):
135 if token_match.group("operation") is not None:
136 return token_match
137 return None
140def _fast_parse_update(
141 update_query: str,
142) -> list[tuple[str, list[tuple[str, str, str, str]]]]:
143 operations: list[tuple[str, list[tuple[str, str, str, str]]]] = []
144 query_len = len(update_query)
145 pos = 0
147 # Each operation is parsed up to the brace that closes it, and the scan
148 # resumes past that brace: a literal quoting "INSERT DATA" would otherwise
149 # read as the start of a new operation and truncate the one it sits in.
150 while pos < query_len:
151 operation_match = _find_next_operation(update_query, pos)
152 if operation_match is None:
153 break
154 block_start = update_query.find("{", operation_match.end())
155 if block_start == -1:
156 break
157 operation_type = (
158 "DeleteData"
159 if operation_match.group("operation").upper() == "DELETE"
160 else "InsertData"
161 )
162 block_end = _find_matching_close_brace(update_query, block_start + 1)
163 operations.append(
164 (
165 operation_type,
166 _parse_graph_blocks(update_query, block_start + 1, block_end),
167 )
168 )
169 pos = block_end + 1
171 return operations
174def _apply_inverse_update(
175 current_state: set[tuple[str, ...]],
176 update_query: str,
177 quad_filter: Callable[[tuple[str, ...]], bool] | None = None,
178) -> None:
179 for operation_type, quads in _fast_parse_update(update_query):
180 if quad_filter is not None:
181 matching_quads = [quad for quad in quads if quad_filter(quad)]
182 else:
183 matching_quads = quads
184 if operation_type == "DeleteData":
185 for quad in matching_quads:
186 current_state.add(quad)
187 elif operation_type == "InsertData":
188 for quad in matching_quads:
189 current_state.discard(quad)
192def _apply_update_ops(
193 operations: list[tuple[str, list[tuple[str, str, str, str]]]],
194 additions: set[tuple[str, ...]],
195 deletions: set[tuple[str, ...]],
196) -> None:
197 for op_type, quads in operations:
198 if op_type == "DeleteData":
199 for quad in quads:
200 if quad in additions:
201 additions.discard(quad)
202 else:
203 deletions.add(quad)
204 elif op_type == "InsertData":
205 for quad in quads:
206 if quad in deletions:
207 deletions.discard(quad)
208 else:
209 additions.add(quad)
212def _compose_update_queries(
213 update_queries: list[str],
214) -> tuple[set[tuple[str, ...]], set[tuple[str, ...]]]:
215 additions: set[tuple[str, ...]] = set()
216 deletions: set[tuple[str, ...]] = set()
217 for uq in update_queries:
218 _apply_update_ops(_fast_parse_update(uq), additions, deletions)
219 return additions, deletions
222def _iter_working_states(
223 sorted_versions: list[tuple[str, str | None]],
224 current_state: set[tuple[str, ...]],
225 target_times: set[str] | None = None,
226 quad_filter: Callable[[tuple[str, ...]], bool] | None = None,
227) -> Iterator[tuple[str, set[tuple[str, ...]]]]:
228 target_count = len(target_times) if target_times is not None else None
229 working_state = set(current_state)
230 materialized_count = 0
231 for index, (timestamp, _update_query) in enumerate(sorted_versions):
232 if index > 0:
233 previous_update = sorted_versions[index - 1][1]
234 if previous_update is not None:
235 _apply_inverse_update(working_state, previous_update, quad_filter)
236 if target_times is None or timestamp in target_times:
237 normalized_timestamp = str(convert_to_datetime(timestamp, stringify=True))
238 yield normalized_timestamp, working_state
239 materialized_count += 1
240 if target_count is not None and materialized_count == target_count:
241 return
244def _materialize_versions(
245 sorted_versions: list[tuple[str, str | None]],
246 current_state: set[tuple[str, ...]],
247 target_times: set[str] | None = None,
248 quad_filter: Callable[[tuple[str, ...]], bool] | None = None,
249) -> list[tuple[str, tuple[tuple[str, ...], ...]]]:
250 return [
251 (timestamp, tuple(working_state))
252 for timestamp, working_state in _iter_working_states(
253 sorted_versions, current_state, target_times, quad_filter
254 )
255 ]
258CONFIG_PATH = "./config.json"
261_GEN_AT_TIME_N3 = f"<{ProvEntity.iri_generated_at_time}>"
262_HAS_UQ_N3 = f"<{ProvEntity.iri_has_update_query}>"
265def _extract_snapshot_update_queries(
266 quads: set[tuple[str, ...]],
267) -> dict[str, str | None]:
268 by_subject: dict[str, dict[str, str]] = {}
269 for quad in quads:
270 if quad[1] in (_GEN_AT_TIME_N3, _HAS_UQ_N3):
271 by_subject.setdefault(quad[0], {})[quad[1]] = _n3_value(quad[2])
272 result: dict[str, str | None] = {}
273 for props in by_subject.values():
274 if _GEN_AT_TIME_N3 in props:
275 result[props[_GEN_AT_TIME_N3]] = props.get(_HAS_UQ_N3)
276 return result
279_PROV_PREFIX = ProvEntity.PROV
280_RDF_TYPE_N3 = f"<{_RDF_TYPE}>"
283def _find_related_object_uris(entity_uri: str, graphs: dict) -> set[str]:
284 entity_n3 = f"<{entity_uri}>"
285 result = set()
286 for quad_set in graphs.values():
287 if quad_set is None:
288 continue
289 for quad in quad_set:
290 if (
291 quad[0] == entity_n3
292 and quad[2].startswith("<")
293 and _PROV_PREFIX not in quad[1]
294 and quad[1] != _RDF_TYPE_N3
295 ):
296 result.add(_n3_value(quad[2]))
297 return result
300class AgnosticEntity:
301 def __init__(
302 self,
303 res: str,
304 config: dict,
305 *,
306 include_related_objects: bool = False,
307 include_merged_entities: bool = False,
308 include_reverse_relations: bool = False,
309 include_historical_reverse_relations: bool = False,
310 reverse_relations_depth: int | None = None,
311 ):
312 self.res = res
313 self.include_related_objects = include_related_objects
314 self.include_merged_entities = include_merged_entities
315 self.include_reverse_relations = include_reverse_relations
316 self.include_historical_reverse_relations = include_historical_reverse_relations
317 self.reverse_relations_depth = reverse_relations_depth
318 self.config = config
320 def get_history(self, *, include_prov_metadata: bool = False) -> tuple:
321 if (
322 self.include_related_objects
323 or self.include_merged_entities
324 or self.include_reverse_relations
325 ):
326 histories = {}
327 self._collect_all_related_entities_histories(
328 histories, include_prov_metadata=include_prov_metadata
329 )
330 return self._get_merged_histories(
331 histories, include_prov_metadata=include_prov_metadata
332 )
333 entity_history = self._get_entity_current_state(
334 include_prov_metadata=include_prov_metadata
335 )
336 entity_history = self._get_old_graphs(entity_history)
337 for uri, time_dict in entity_history[0].items():
338 for ts, quad_set in time_dict.items():
339 if quad_set is None:
340 entity_history[0][uri][ts] = set()
341 return tuple(entity_history)
343 def get_histories_by_entity(self, *, include_prov_metadata: bool = False) -> tuple:
344 histories = {}
345 self._collect_all_related_entities_histories(
346 histories, include_prov_metadata=include_prov_metadata
347 )
348 entity_histories = {}
349 metadata = {}
350 for entity_uri, (entity_history_dict, entity_metadata) in histories.items():
351 entity_states = entity_history_dict[entity_uri]
352 if not entity_states:
353 continue
354 entity_histories[entity_uri] = {
355 timestamp: quad_set if quad_set is not None else set()
356 for timestamp, quad_set in entity_states.items()
357 }
358 if include_prov_metadata and entity_metadata:
359 metadata[entity_uri] = entity_metadata[entity_uri]
360 return entity_histories, metadata
362 def _collect_all_related_entities_histories(
363 self, histories: dict, *, include_prov_metadata: bool
364 ) -> None:
365 main_entity = AgnosticEntity(
366 self.res,
367 self.config,
368 include_related_objects=False,
369 include_merged_entities=False,
370 include_reverse_relations=False,
371 )
372 entity_history = main_entity._get_entity_current_state(
373 include_prov_metadata=include_prov_metadata
374 )
375 entity_history = main_entity._get_old_graphs(entity_history)
376 histories[self.res] = (entity_history[0], entity_history[1])
378 processed_entities = {self.res}
380 if self.include_related_objects:
381 self._collect_related_objects_recursively(
382 self.res,
383 processed_entities,
384 histories,
385 include_prov_metadata=include_prov_metadata,
386 )
388 if self.include_merged_entities:
389 self._collect_merged_entities_recursively(
390 self.res,
391 processed_entities,
392 histories,
393 include_prov_metadata=include_prov_metadata,
394 )
396 if self.include_reverse_relations:
397 self._collect_reverse_relations_recursively(
398 self.res,
399 processed_entities,
400 histories,
401 include_prov_metadata=include_prov_metadata,
402 depth=self.reverse_relations_depth,
403 )
405 def _collect_related_objects_recursively(
406 self,
407 entity_uri: str,
408 processed_entities: set[str],
409 histories: dict,
410 *,
411 include_prov_metadata: bool,
412 depth: int | None = None,
413 ) -> None:
414 if depth is not None and depth <= 0:
415 return
417 next_depth = None if depth is None else depth - 1
419 entity_graphs = (
420 histories[entity_uri][0][entity_uri] if entity_uri in histories else None
421 )
422 if not entity_graphs:
423 return
425 for obj_uri in _find_related_object_uris(entity_uri, entity_graphs):
426 if obj_uri not in processed_entities:
427 processed_entities.add(obj_uri)
428 agnostic_entity = AgnosticEntity(
429 obj_uri,
430 self.config,
431 include_related_objects=False,
432 include_merged_entities=False,
433 include_reverse_relations=False,
434 )
435 entity_history = agnostic_entity._get_entity_current_state(
436 include_prov_metadata=include_prov_metadata
437 )
438 entity_history = agnostic_entity._get_old_graphs(entity_history)
439 histories[obj_uri] = (entity_history[0], entity_history[1])
440 self._collect_related_objects_recursively(
441 obj_uri,
442 processed_entities,
443 histories,
444 include_prov_metadata=include_prov_metadata,
445 depth=next_depth,
446 )
448 def _collect_merged_entities_recursively(
449 self,
450 entity_uri: str,
451 processed_entities: set[str],
452 histories: dict,
453 *,
454 include_prov_metadata: bool,
455 depth: int | None = None,
456 ) -> None:
457 if depth is not None and depth <= 0:
458 return
460 next_depth = None if depth is None else depth - 1
462 merged_entities = self._find_merged_entities(entity_uri)
464 for merged_entity_uri in merged_entities:
465 if merged_entity_uri not in processed_entities:
466 processed_entities.add(merged_entity_uri)
467 agnostic_entity = AgnosticEntity(
468 merged_entity_uri,
469 self.config,
470 include_related_objects=False,
471 include_merged_entities=False,
472 include_reverse_relations=False,
473 )
474 entity_history = agnostic_entity._get_entity_current_state(
475 include_prov_metadata=include_prov_metadata
476 )
477 entity_history = agnostic_entity._get_old_graphs(entity_history)
478 histories[merged_entity_uri] = (entity_history[0], entity_history[1])
479 self._collect_merged_entities_recursively(
480 merged_entity_uri,
481 processed_entities,
482 histories,
483 include_prov_metadata=include_prov_metadata,
484 depth=next_depth,
485 )
487 def _collect_reverse_relations_recursively(
488 self,
489 entity_uri: str,
490 processed_entities: set[str],
491 histories: dict,
492 *,
493 include_prov_metadata: bool,
494 depth: int | None = None,
495 ) -> None:
496 if depth is not None and depth <= 0:
497 return
499 next_depth = None if depth is None else depth - 1
501 reverse_related_entities = self._find_reverse_related_entities(entity_uri)
503 for reverse_entity_uri in reverse_related_entities:
504 if reverse_entity_uri not in processed_entities:
505 processed_entities.add(reverse_entity_uri)
506 agnostic_entity = AgnosticEntity(
507 reverse_entity_uri,
508 self.config,
509 include_related_objects=False,
510 include_merged_entities=False,
511 include_reverse_relations=False,
512 )
513 entity_history = agnostic_entity._get_entity_current_state(
514 include_prov_metadata=include_prov_metadata
515 )
516 entity_history = agnostic_entity._get_old_graphs(entity_history)
517 histories[reverse_entity_uri] = (entity_history[0], entity_history[1])
518 self._collect_reverse_relations_recursively(
519 reverse_entity_uri,
520 processed_entities,
521 histories,
522 include_prov_metadata=include_prov_metadata,
523 depth=next_depth,
524 )
526 def _get_merged_histories(
527 self, histories: dict, *, include_prov_metadata: bool
528 ) -> tuple:
529 entity_histories = {}
530 metadata = {}
531 for entity_uri, (entity_history_dict, entity_metadata) in histories.items():
532 entity_histories[entity_uri] = entity_history_dict[entity_uri]
533 if include_prov_metadata and entity_metadata:
534 metadata[entity_uri] = entity_metadata[entity_uri]
536 main_entity_times = sorted(
537 entity_histories[self.res].keys(), key=_parse_datetime
538 )
540 merged_histories = {self.res: {}}
542 related_sorted_times = {}
543 for entity_uri, entity_history in entity_histories.items():
544 if entity_uri == self.res:
545 continue
546 related_sorted_times[entity_uri] = sorted(
547 ((t, _parse_datetime(t)) for t in entity_history), key=lambda x: x[1]
548 )
550 for timestamp in main_entity_times:
551 merged_set = set(entity_histories[self.res][timestamp])
552 timestamp_dt = _parse_datetime(timestamp)
554 for entity_uri, sorted_times in related_sorted_times.items():
555 relevant_time = None
556 for etime, etime_dt in sorted_times:
557 if etime_dt <= timestamp_dt:
558 relevant_time = etime
559 else:
560 break
561 if relevant_time:
562 merged_set.update(entity_histories[entity_uri][relevant_time])
564 merged_histories[self.res][timestamp] = merged_set
566 return merged_histories, metadata
568 def get_state_at_time(
569 self,
570 time: tuple[str | None, str | None],
571 *,
572 include_prov_metadata: bool = False,
573 ) -> tuple:
574 if (
575 self.include_related_objects
576 or self.include_merged_entities
577 or self.include_reverse_relations
578 ):
579 histories = {}
580 self._collect_all_related_entities_states_at_time(
581 histories, time, include_prov_metadata=include_prov_metadata
582 )
583 return self._get_merged_histories_at_time(
584 histories, include_prov_metadata=include_prov_metadata
585 )
586 return self._get_entity_state_at_time(
587 time, include_prov_metadata=include_prov_metadata
588 )
590 def get_delta(
591 self,
592 time_start: str,
593 time_end: str,
594 ) -> tuple[set[tuple[str, ...]], set[tuple[str, ...]]]:
595 is_quadstore = self.config["provenance"]["is_quadstore"]
596 graph_statement = f"GRAPH <{self.res}/prov/>" if is_quadstore else ""
597 query_snapshots = f"""
598 SELECT ?time ?updateQuery
599 WHERE {{
600 {graph_statement}
601 {{
602 ?snapshot <{ProvEntity.iri_specialization_of}> <{self.res}>;
603 <{ProvEntity.iri_generated_at_time}> ?time.
604 OPTIONAL {{
605 ?snapshot <{ProvEntity.iri_has_update_query}> ?updateQuery.
606 }}
607 }}
608 }}
609 """
610 results = Sparql(query_snapshots, config=self.config).run_select_query()
611 bindings = results["results"]["bindings"]
612 if not bindings:
613 return set(), set()
614 start_dt = _parse_datetime(time_start)
615 end_dt = _parse_datetime(time_end)
616 parsed = [(b, _parse_datetime(b["time"]["value"])) for b in bindings]
617 first_snapshot_dt = min(dt for _, dt in parsed)
618 if first_snapshot_dt > start_dt:
619 entity_graphs, _, _ = self._get_entity_state_at_time(
620 (time_end, time_end), include_prov_metadata=False
621 )
622 if not entity_graphs:
623 return set(), set()
624 state_at_end = next(iter(entity_graphs.values()))
625 return state_at_end, set()
626 relevant = sorted(
627 (
628 (b, dt)
629 for b, dt in parsed
630 if start_dt < dt <= end_dt
631 and "updateQuery" in b
632 and "value" in b["updateQuery"]
633 ),
634 key=lambda x: x[1],
635 )
636 return _compose_update_queries([b["updateQuery"]["value"] for b, _ in relevant])
638 def _collect_all_related_entities_states_at_time(
639 self,
640 histories: dict,
641 time: tuple[str | None, str | None],
642 *,
643 include_prov_metadata: bool,
644 ) -> None:
645 main_entity = AgnosticEntity(
646 self.res,
647 self.config,
648 include_related_objects=False,
649 include_merged_entities=False,
650 include_reverse_relations=False,
651 )
652 entity_graphs, entity_snapshots, other_snapshots_metadata = (
653 main_entity._get_entity_state_at_time(
654 time, include_prov_metadata=include_prov_metadata
655 )
656 )
657 histories[self.res] = (
658 entity_graphs,
659 entity_snapshots,
660 other_snapshots_metadata,
661 )
663 processed_entities = {self.res}
665 if self.include_related_objects:
666 self._collect_related_objects_states_at_time(
667 self.res,
668 processed_entities,
669 histories,
670 time,
671 include_prov_metadata=include_prov_metadata,
672 )
674 if self.include_merged_entities:
675 self._collect_merged_entities_states_at_time(
676 self.res,
677 processed_entities,
678 histories,
679 time,
680 include_prov_metadata=include_prov_metadata,
681 )
683 if self.include_reverse_relations:
684 self._collect_reverse_relations_states_at_time(
685 self.res,
686 processed_entities,
687 histories,
688 time,
689 include_prov_metadata=include_prov_metadata,
690 depth=self.reverse_relations_depth,
691 )
693 def _collect_related_objects_states_at_time(
694 self,
695 entity_uri: str,
696 processed_entities: set[str],
697 histories: dict,
698 time: tuple[str | None, str | None],
699 *,
700 include_prov_metadata: bool,
701 depth: int | None = None,
702 ) -> None:
703 if depth is not None and depth <= 0:
704 return
706 next_depth = None if depth is None else depth - 1
708 entity_graphs = histories[entity_uri][0] if entity_uri in histories else None
709 if not entity_graphs:
710 return
712 for obj_uri in _find_related_object_uris(entity_uri, entity_graphs):
713 if obj_uri not in processed_entities:
714 processed_entities.add(obj_uri)
715 agnostic_entity = AgnosticEntity(
716 obj_uri,
717 self.config,
718 include_related_objects=False,
719 include_merged_entities=False,
720 include_reverse_relations=False,
721 )
722 entity_graphs_new, entity_snapshots, other_snapshots_metadata = (
723 agnostic_entity._get_entity_state_at_time(
724 time, include_prov_metadata=include_prov_metadata
725 )
726 )
727 histories[obj_uri] = (
728 entity_graphs_new,
729 entity_snapshots,
730 other_snapshots_metadata,
731 )
732 self._collect_related_objects_states_at_time(
733 obj_uri,
734 processed_entities,
735 histories,
736 time,
737 include_prov_metadata=include_prov_metadata,
738 depth=next_depth,
739 )
741 def _collect_merged_entities_states_at_time(
742 self,
743 entity_uri: str,
744 processed_entities: set[str],
745 histories: dict,
746 time: tuple[str | None, str | None],
747 *,
748 include_prov_metadata: bool,
749 depth: int | None = None,
750 ) -> None:
751 if depth is not None and depth <= 0:
752 return
754 next_depth = None if depth is None else depth - 1
756 merged_entities = self._find_merged_entities(entity_uri)
758 for merged_entity_uri in merged_entities:
759 if merged_entity_uri not in processed_entities:
760 processed_entities.add(merged_entity_uri)
761 agnostic_entity = AgnosticEntity(
762 merged_entity_uri,
763 self.config,
764 include_related_objects=False,
765 include_merged_entities=False,
766 include_reverse_relations=False,
767 )
768 entity_graphs, entity_snapshots, other_snapshots_metadata = (
769 agnostic_entity._get_entity_state_at_time(
770 time, include_prov_metadata=include_prov_metadata
771 )
772 )
773 histories[merged_entity_uri] = (
774 entity_graphs,
775 entity_snapshots,
776 other_snapshots_metadata,
777 )
778 self._collect_merged_entities_states_at_time(
779 merged_entity_uri,
780 processed_entities,
781 histories,
782 time,
783 include_prov_metadata=include_prov_metadata,
784 depth=next_depth,
785 )
787 def _collect_reverse_relations_states_at_time(
788 self,
789 entity_uri: str,
790 processed_entities: set[str],
791 histories: dict,
792 time: tuple[str | None, str | None],
793 *,
794 include_prov_metadata: bool,
795 depth: int | None = None,
796 ) -> None:
797 if depth is not None and depth <= 0:
798 return
800 next_depth = None if depth is None else depth - 1
802 reverse_related_entities = self._find_reverse_related_entities(entity_uri)
804 for reverse_entity_uri in reverse_related_entities:
805 if reverse_entity_uri not in processed_entities:
806 processed_entities.add(reverse_entity_uri)
807 agnostic_entity = AgnosticEntity(
808 reverse_entity_uri,
809 self.config,
810 include_related_objects=False,
811 include_merged_entities=False,
812 include_reverse_relations=False,
813 )
814 entity_graphs, entity_snapshots, other_snapshots_metadata = (
815 agnostic_entity._get_entity_state_at_time(
816 time, include_prov_metadata=include_prov_metadata
817 )
818 )
819 histories[reverse_entity_uri] = (
820 entity_graphs,
821 entity_snapshots,
822 other_snapshots_metadata,
823 )
824 self._collect_reverse_relations_states_at_time(
825 reverse_entity_uri,
826 processed_entities,
827 histories,
828 time,
829 include_prov_metadata=include_prov_metadata,
830 depth=next_depth,
831 )
833 def _get_merged_histories_at_time(
834 self, histories: dict, *, include_prov_metadata: bool
835 ) -> tuple:
836 entity_histories = {}
837 entity_snapshots_metadata = {}
838 other_snapshots_metadata = {} if include_prov_metadata else None
840 for entity_uri, (
841 entity_graphs,
842 entity_snapshots,
843 other_snapshots,
844 ) in histories.items():
845 entity_histories[entity_uri] = entity_graphs
846 entity_snapshots_metadata[entity_uri] = entity_snapshots
847 if (
848 include_prov_metadata
849 and other_snapshots
850 and other_snapshots_metadata is not None
851 ):
852 other_snapshots_metadata[entity_uri] = other_snapshots
854 main_entity_times = sorted(
855 set(entity_histories[self.res].keys()), key=_parse_datetime
856 )
858 merged_histories = {self.res: {}}
860 related_sorted_times = {}
861 for entity_uri, graphs_at_times in entity_histories.items():
862 if entity_uri == self.res:
863 continue
864 related_sorted_times[entity_uri] = sorted(
865 ((t, _parse_datetime(t)) for t in graphs_at_times), key=lambda x: x[1]
866 )
868 for timestamp in main_entity_times:
869 merged_set = set(entity_histories[self.res][timestamp])
870 timestamp_dt = _parse_datetime(timestamp)
872 for entity_uri, sorted_times in related_sorted_times.items():
873 graphs_at_times = entity_histories[entity_uri]
874 if timestamp in graphs_at_times:
875 related_quads = graphs_at_times[timestamp]
876 else:
877 relevant_time = None
878 for rt, rt_dt in sorted_times:
879 if rt_dt <= timestamp_dt:
880 relevant_time = rt
881 else:
882 break
883 if relevant_time:
884 related_quads = graphs_at_times[relevant_time]
885 else:
886 continue
887 merged_set.update(related_quads)
889 merged_histories[self.res][timestamp] = merged_set
891 return merged_histories, entity_snapshots_metadata, other_snapshots_metadata
893 def _get_entity_state_at_time(
894 self, time: tuple[str | None, str | None], *, include_prov_metadata: bool
895 ) -> tuple:
896 other_snapshots_metadata = {}
897 is_quadstore = self.config["provenance"]["is_quadstore"]
898 graph_statement = f"GRAPH <{self.res}/prov/>" if is_quadstore else ""
899 if include_prov_metadata:
900 query_snapshots = f"""
901 SELECT ?snapshot ?time ?responsibleAgent ?updateQuery
902 ?primarySource ?description ?invalidatedAtTime ?derivedFrom
903 WHERE {{
904 {graph_statement}
905 {{
906 ?snapshot <{ProvEntity.iri_specialization_of}> <{self.res}>;
907 <{ProvEntity.iri_generated_at_time}> ?time;
908 <{ProvEntity.iri_was_attributed_to}> ?responsibleAgent.
909 OPTIONAL {{
910 ?snapshot <{ProvEntity.iri_invalidated_at_time}>
911 ?invalidatedAtTime.
912 }}
913 OPTIONAL {{
914 ?snapshot <{ProvEntity.iri_description}> ?description.
915 }}
916 OPTIONAL {{
917 ?snapshot <{ProvEntity.iri_has_update_query}> ?updateQuery.
918 }}
919 OPTIONAL {{
920 ?snapshot <{ProvEntity.iri_had_primary_source}>
921 ?primarySource.
922 }}
923 OPTIONAL {{
924 ?snapshot <{ProvEntity.iri_was_derived_from}> ?derivedFrom.
925 }}
926 }}
927 }}
928 """
929 else:
930 query_snapshots = f"""
931 SELECT ?snapshot ?time ?updateQuery
932 WHERE {{
933 {graph_statement}
934 {{
935 ?snapshot <{ProvEntity.iri_specialization_of}> <{self.res}>;
936 <{ProvEntity.iri_generated_at_time}> ?time.
937 OPTIONAL {{
938 ?snapshot <{ProvEntity.iri_has_update_query}> ?updateQuery.
939 }}
940 }}
941 }}
942 """
943 results = Sparql(query_snapshots, config=self.config).run_select_query()
944 bindings = results["results"]["bindings"]
945 if not bindings:
946 return {}, {}, other_snapshots_metadata
947 snapshots_by_uri = {}
948 for binding in bindings:
949 snapshots_by_uri.setdefault(binding["snapshot"]["value"], binding)
950 sorted_results = sorted(
951 snapshots_by_uri.values(),
952 key=lambda x: _parse_datetime(x["time"]["value"]),
953 reverse=True,
954 )
955 relevant_results, start_timestamp_alias = _select_interval_snapshots(
956 time, sorted_results, time_index="time"
957 )
958 relevant_snapshot_uris = {
959 result["snapshot"]["value"] for result in relevant_results
960 }
961 entity_snapshots = {}
962 if include_prov_metadata:
963 metadata_by_snapshot = {}
964 for result in bindings:
965 snapshot_uri = result["snapshot"]["value"]
966 snapshot_metadata = metadata_by_snapshot.setdefault(
967 snapshot_uri,
968 {
969 "generatedAtTime": result["time"]["value"],
970 "invalidatedAtTime": result.get("invalidatedAtTime", {}).get(
971 "value"
972 ),
973 "wasAttributedTo": result["responsibleAgent"]["value"],
974 "hasUpdateQuery": result.get("updateQuery", {}).get("value"),
975 "hadPrimarySource": result.get("primarySource", {}).get(
976 "value"
977 ),
978 "description": result.get("description", {}).get("value"),
979 "wasDerivedFrom": [],
980 },
981 )
982 if "derivedFrom" in result:
983 snapshot_metadata["wasDerivedFrom"].append(
984 result["derivedFrom"]["value"]
985 )
986 for snapshot_metadata in metadata_by_snapshot.values():
987 snapshot_metadata["wasDerivedFrom"] = sorted(
988 set(snapshot_metadata["wasDerivedFrom"])
989 )
990 entity_snapshots = {
991 snapshot_uri: metadata
992 for snapshot_uri, metadata in metadata_by_snapshot.items()
993 if snapshot_uri in relevant_snapshot_uris
994 }
995 other_snapshots_metadata = {
996 snapshot_uri: metadata
997 for snapshot_uri, metadata in metadata_by_snapshot.items()
998 if snapshot_uri not in relevant_snapshot_uris
999 }
1000 if not relevant_results:
1001 return {}, entity_snapshots, other_snapshots_metadata
1002 entity_quads = self._query_dataset(self.res)
1003 sorted_versions = [
1004 (
1005 result["time"]["value"],
1006 result["updateQuery"]["value"]
1007 if "updateQuery" in result and "value" in result["updateQuery"]
1008 else None,
1009 )
1010 for result in sorted_results
1011 ]
1012 target_times = {
1013 relevant_result["time"]["value"] for relevant_result in relevant_results
1014 }
1015 entity_graphs = {}
1016 for timestamp, quad_set in _materialize_versions(
1017 sorted_versions, entity_quads, target_times
1018 ):
1019 result_timestamp = (
1020 start_timestamp_alias[1]
1021 if start_timestamp_alias is not None
1022 and timestamp == start_timestamp_alias[0]
1023 else timestamp
1024 )
1025 entity_graphs[result_timestamp] = set(quad_set)
1026 return entity_graphs, entity_snapshots, other_snapshots_metadata
1028 def _include_prov_metadata(
1029 self, triples_generated_at_time: list, current_state: set[tuple[str, ...]]
1030 ) -> dict:
1031 res_n3 = f"<{self.res}>"
1032 entity_n3 = f"<{ProvEntity.iri_entity}>"
1033 for quad in current_state:
1034 if quad[0] == res_n3 and quad[1] == _RDF_TYPE_N3 and quad[2] == entity_n3:
1035 return {}
1036 prov_properties = {
1037 f"<{ProvEntity.iri_invalidated_at_time}>": "invalidatedAtTime",
1038 f"<{ProvEntity.iri_was_attributed_to}>": "wasAttributedTo",
1039 f"<{ProvEntity.iri_had_primary_source}>": "hadPrimarySource",
1040 f"<{ProvEntity.iri_description}>": "description",
1041 f"<{ProvEntity.iri_has_update_query}>": "hasUpdateQuery",
1042 f"<{ProvEntity.iri_was_derived_from}>": "wasDerivedFrom",
1043 }
1044 prov_metadata: dict = {self.res: {}}
1045 for triple in triples_generated_at_time:
1046 time = convert_to_datetime(_n3_value(triple[2]), stringify=True)
1047 snapshot_uri_str = _n3_value(triple[0])
1048 prov_metadata[self.res][snapshot_uri_str] = {
1049 "generatedAtTime": time,
1050 "invalidatedAtTime": None,
1051 "wasAttributedTo": None,
1052 "hadPrimarySource": None,
1053 "description": None,
1054 "hasUpdateQuery": None,
1055 "wasDerivedFrom": [],
1056 }
1057 prov_prop_n3_set = set(prov_properties)
1058 index: dict[str, dict[str, list[str]]] = {}
1059 for quad in current_state:
1060 if quad[1] in prov_prop_n3_set:
1061 index.setdefault(quad[0], {}).setdefault(quad[1], []).append(
1062 _n3_value(quad[2])
1063 )
1064 for metadata in dict(prov_metadata).values():
1065 for se_uri_str, snapshot_data in metadata.items():
1066 se_n3 = f"<{se_uri_str}>"
1067 se_props = index.get(se_n3, {})
1068 for prov_prop_n3, abbr in prov_properties.items():
1069 for value in se_props.get(prov_prop_n3, ()):
1070 if abbr == "wasDerivedFrom":
1071 snapshot_data[abbr].append(value)
1072 else:
1073 snapshot_data[abbr] = value
1074 if isinstance(snapshot_data.get("wasDerivedFrom"), list):
1075 snapshot_data["wasDerivedFrom"] = sorted(
1076 snapshot_data["wasDerivedFrom"]
1077 )
1079 return prov_metadata
1081 def _get_entity_current_state(self, *, include_prov_metadata: bool = False) -> list:
1082 entity_current_state: list = [{self.res: {}}]
1083 prov_quads = self._query_provenance(include_prov_metadata=include_prov_metadata)
1084 if len(prov_quads) == 0:
1085 entity_current_state.append({})
1086 return entity_current_state
1087 dataset_quads = self._query_dataset(self.res)
1088 gen_at_time_n3 = f"<{ProvEntity.iri_generated_at_time}>"
1089 triples_generated_at_time = [
1090 quad for quad in prov_quads if quad[1] == gen_at_time_n3
1091 ]
1092 most_recent_time = None
1093 most_recent_time_str: str | None = None
1094 for quad in triples_generated_at_time:
1095 snapshot_time_str = _n3_value(quad[2])
1096 snapshot_date_time = _parse_datetime(snapshot_time_str)
1097 if most_recent_time:
1098 if snapshot_date_time > most_recent_time:
1099 most_recent_time = snapshot_date_time
1100 most_recent_time_str = snapshot_time_str
1101 else:
1102 most_recent_time = snapshot_date_time
1103 most_recent_time_str = snapshot_time_str
1104 entity_current_state[0][self.res][snapshot_time_str] = None
1105 entity_current_state[0][self.res][most_recent_time_str] = dataset_quads
1106 if include_prov_metadata:
1107 prov_metadata = self._include_prov_metadata(
1108 triples_generated_at_time, prov_quads
1109 )
1110 entity_current_state.append(prov_metadata)
1111 else:
1112 entity_current_state.append(None)
1113 entity_current_state.append(prov_quads)
1114 return entity_current_state
1116 def _get_old_graphs(self, entity_current_state: list) -> list:
1117 prov_quads_index = 2
1118 prov_quads = (
1119 entity_current_state.pop(prov_quads_index)
1120 if len(entity_current_state) > prov_quads_index
1121 else set()
1122 )
1123 snapshot_update_queries = _extract_snapshot_update_queries(prov_quads)
1124 ordered_data: list[tuple[str, set[tuple[str, ...]]]] = sorted(
1125 entity_current_state[0][self.res].items(),
1126 key=lambda x: _parse_datetime(str(x[0])),
1127 reverse=True,
1128 )
1129 if not ordered_data:
1130 return entity_current_state
1131 for index, date_graph in enumerate(ordered_data):
1132 if index > 0:
1133 next_snapshot = ordered_data[index - 1][0]
1134 previous_graph = set(entity_current_state[0][self.res][next_snapshot])
1135 update_query = snapshot_update_queries.get(str(next_snapshot))
1136 if update_query is None:
1137 entity_current_state[0][self.res][date_graph[0]] = previous_graph
1138 else:
1139 _apply_inverse_update(previous_graph, update_query)
1140 entity_current_state[0][self.res][date_graph[0]] = previous_graph
1141 for time in list(entity_current_state[0][self.res]):
1142 quad_set = entity_current_state[0][self.res].pop(time)
1143 time_str = str(convert_to_datetime(str(time), stringify=True))
1144 entity_current_state[0][self.res][time_str] = quad_set
1145 return entity_current_state
1147 def iter_versions(self):
1148 prov_quads = self._query_provenance(include_prov_metadata=False)
1149 if len(prov_quads) == 0:
1150 return
1151 dataset_quads = self._query_dataset(self.res)
1152 working: set[tuple[str, ...]] = set(dataset_quads)
1153 snapshots = _extract_snapshot_update_queries(prov_quads)
1154 ordered = sorted(
1155 snapshots.items(), key=lambda x: _parse_datetime(x[0]), reverse=True
1156 )
1157 for i, (time_str, _update_query) in enumerate(ordered):
1158 if i > 0:
1159 prev_update = ordered[i - 1][1]
1160 if prev_update is not None:
1161 _apply_inverse_update(working, prev_update)
1162 normalized = str(convert_to_datetime(time_str, stringify=True))
1163 yield normalized, set(working)
1165 def _query_dataset(self, entity_uri: str | None = None) -> set[tuple[str, ...]]:
1166 entity_uri = self.res if entity_uri is None else entity_uri
1168 is_quadstore = self.config["dataset"]["is_quadstore"]
1170 if is_quadstore:
1171 query_dataset = f"""
1172 SELECT ?s ?p ?o ?g
1173 WHERE {{
1174 GRAPH ?g {{
1175 VALUES ?s {{<{entity_uri}>}}
1176 ?s ?p ?o
1177 }}
1178 }}
1179 """
1180 else:
1181 query_dataset = f"""
1182 SELECT ?s ?p ?o
1183 WHERE {{
1184 VALUES ?s {{<{entity_uri}>}}
1185 ?s ?p ?o
1186 }}
1187 """
1189 return Sparql(query_dataset, config=self.config).run_select_to_quad_set()
1191 def _query_provenance(
1192 self, *, include_prov_metadata: bool = False
1193 ) -> set[tuple[str, ...]]:
1194 if include_prov_metadata:
1195 query_provenance = f"""
1196 SELECT ?s ?p ?o WHERE {{
1197 ?s <{ProvEntity.iri_specialization_of}> <{self.res}>;
1198 <{ProvEntity.iri_was_attributed_to}> ?_agent;
1199 <{ProvEntity.iri_generated_at_time}> ?_t;
1200 <{ProvEntity.iri_description}> ?_desc.
1201 ?s ?p ?o.
1202 VALUES ?p {{
1203 <{ProvEntity.iri_generated_at_time}>
1204 <{ProvEntity.iri_was_attributed_to}>
1205 <{ProvEntity.iri_had_primary_source}>
1206 <{ProvEntity.iri_description}>
1207 <{ProvEntity.iri_has_update_query}>
1208 <{ProvEntity.iri_invalidated_at_time}>
1209 <{ProvEntity.iri_was_derived_from}>
1210 <{ProvEntity.iri_specialization_of}>
1211 }}
1212 }}
1213 """
1214 else:
1215 query_provenance = f"""
1216 SELECT ?s ?p ?o WHERE {{
1217 ?s <{ProvEntity.iri_specialization_of}> <{self.res}>;
1218 <{ProvEntity.iri_generated_at_time}> ?_t.
1219 ?s ?p ?o.
1220 VALUES ?p {{
1221 <{ProvEntity.iri_generated_at_time}>
1222 <{ProvEntity.iri_has_update_query}>
1223 <{ProvEntity.iri_was_derived_from}>
1224 <{ProvEntity.iri_specialization_of}>
1225 }}
1226 }}
1227 """
1228 return Sparql(query_provenance, config=self.config).run_select_to_quad_set()
1230 def _find_merged_entities(self, entity_uri: str) -> set[str]:
1231 merged_entity_uris = set()
1232 query_simple = f"""
1233 SELECT ?merged_entity_uri
1234 WHERE {{
1235 ?snapshot <{ProvEntity.iri_specialization_of}> <{entity_uri}> .
1236 ?snapshot <{ProvEntity.iri_was_derived_from}> ?derived_snapshot .
1237 ?derived_snapshot <{ProvEntity.iri_specialization_of}>
1238 ?merged_entity_uri .
1239 FILTER (?merged_entity_uri != <{entity_uri}>)
1240 }}
1241 """
1242 results = Sparql(query_simple, config=self.config).run_select_query()
1243 bindings = results.get("results", {}).get("bindings", [])
1244 for binding in bindings:
1245 if (
1246 "merged_entity_uri" in binding
1247 and "value" in binding["merged_entity_uri"]
1248 ):
1249 merged_entity_uris.add(binding["merged_entity_uri"]["value"])
1251 return merged_entity_uris
1253 def _find_reverse_related_entities(self, entity_uri: str) -> set[str]:
1254 reverse_related_entity_uris = set()
1256 is_quadstore = self.config["dataset"]["is_quadstore"]
1258 if is_quadstore:
1259 query = f"""
1260 SELECT ?subject
1261 WHERE {{
1262 GRAPH ?g {{
1263 ?subject ?predicate <{entity_uri}> .
1264 FILTER(?predicate != <{_RDF_TYPE}>
1265 && !strstarts(str(?predicate), "{ProvEntity.PROV}"))
1266 }}
1267 }}
1268 """
1269 else:
1270 query = f"""
1271 SELECT ?subject
1272 WHERE {{
1273 ?subject ?predicate <{entity_uri}> .
1274 FILTER(?predicate != <{_RDF_TYPE}>
1275 && !strstarts(str(?predicate), "{ProvEntity.PROV}"))
1276 }}
1277 """
1279 self._add_reverse_related_entities_from_query(
1280 query, entity_uri, reverse_related_entity_uris
1281 )
1283 # No index can serve CONTAINS on a literal, so this scans every update
1284 # query in the provenance store. It is the only way to reach entities
1285 # whose reference was deleted, but it costs a full scan per call and is
1286 # therefore requested explicitly rather than by default.
1287 if self.include_historical_reverse_relations:
1288 prov_query = f"""
1289 SELECT DISTINCT ?subject
1290 WHERE {{
1291 ?snapshot <{ProvEntity.iri_specialization_of}> ?subject .
1292 ?snapshot <{ProvEntity.iri_has_update_query}> ?update_query .
1293 FILTER(CONTAINS(?update_query, "<{entity_uri}>"))
1294 }}
1295 """
1296 self._add_reverse_related_entities_from_query(
1297 prov_query, entity_uri, reverse_related_entity_uris
1298 )
1300 return reverse_related_entity_uris
1302 def _add_reverse_related_entities_from_query(
1303 self, select_query: str, entity_uri: str, reverse_related_entity_uris: set[str]
1304 ) -> None:
1305 results = Sparql(select_query, config=self.config).run_select_query()
1306 bindings = results.get("results", {}).get("bindings", [])
1307 for binding in bindings:
1308 if "subject" in binding and "value" in binding["subject"]:
1309 subject_uri = binding["subject"]["value"]
1310 if subject_uri != entity_uri:
1311 reverse_related_entity_uris.add(subject_uri)
1314def _filter_timestamps_by_interval(
1315 interval: tuple[str | None, str | None] | None,
1316 iterator: list,
1317 time_index: str | None = None,
1318) -> list:
1319 if interval:
1320 after_time = _parse_datetime(interval[0]) if interval[0] else None
1321 before_time = _parse_datetime(interval[1]) if interval[1] else None
1322 relevant_timestamps = []
1323 for timestamp in iterator:
1324 if time_index is not None and time_index in timestamp:
1325 time_binding = timestamp[time_index]
1326 if "value" in time_binding:
1327 time_str = time_binding["value"]
1328 time = _parse_datetime(time_str)
1329 else:
1330 continue
1331 else:
1332 continue
1333 if after_time and before_time:
1334 if after_time <= time <= before_time:
1335 relevant_timestamps.append(timestamp)
1336 elif after_time and not before_time:
1337 if time >= after_time:
1338 relevant_timestamps.append(timestamp)
1339 elif before_time and not after_time:
1340 if time <= before_time:
1341 relevant_timestamps.append(timestamp)
1342 else:
1343 relevant_timestamps.append(timestamp)
1344 else:
1345 relevant_timestamps = iterator.copy()
1346 return relevant_timestamps
1349def _select_interval_snapshots(
1350 interval: tuple[str | None, str | None] | None,
1351 snapshots: list[dict],
1352 *,
1353 time_index: str,
1354) -> tuple[list[dict], tuple[str, str] | None]:
1355 selected_snapshots = _filter_timestamps_by_interval(
1356 interval, snapshots, time_index=time_index
1357 )
1358 if not interval or not interval[0]:
1359 return selected_snapshots, None
1360 interval_start = _parse_datetime(interval[0])
1361 snapshots_at_start = [
1362 snapshot
1363 for snapshot in snapshots
1364 if _parse_datetime(snapshot[time_index]["value"]) <= interval_start
1365 ]
1366 if not snapshots_at_start:
1367 return selected_snapshots, None
1368 start_snapshot = max(
1369 snapshots_at_start,
1370 key=lambda snapshot: _parse_datetime(snapshot[time_index]["value"]),
1371 )
1372 if start_snapshot not in selected_snapshots:
1373 selected_snapshots.append(start_snapshot)
1374 timestamp_alias = (
1375 str(convert_to_datetime(start_snapshot[time_index]["value"], stringify=True)),
1376 str(convert_to_datetime(interval[0], stringify=True)),
1377 )
1378 return selected_snapshots, timestamp_alias