Coverage for oc_meta / core / curator.py: 95%
962 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# SPDX-FileCopyrightText: 2019 Silvio Peroni <silvio.peroni@unibo.it>
2# SPDX-FileCopyrightText: 2019-2020 Fabio Mariani <fabio.mariani555@gmail.com>
3# SPDX-FileCopyrightText: 2021 Simone Persiani <iosonopersia@gmail.com>
4# SPDX-FileCopyrightText: 2021-2026 Arcangelo Massari <arcangelo.massari@unibo.it>
5#
6# SPDX-License-Identifier: ISC
8from __future__ import annotations
10import multiprocessing
11import os
12from concurrent.futures import ProcessPoolExecutor
13from contextlib import nullcontext
14from typing import TYPE_CHECKING, Dict, List, Tuple
16from oc_meta.constants import CONTAINER_EDITOR_TYPES, VALID_ENTITY_TYPES
17from oc_meta.lib.cleaner import (
18 clean_date,
19 clean_name,
20 clean_ra_list,
21 clean_title,
22 clean_volume_and_issue,
23 normalize_hyphens,
24 normalize_id,
25)
26from oc_meta.lib.file_manager import write_csv
27from oc_meta.lib.finder import ResourceFinder
28from oc_meta.lib.merge_registry import EntityStore
29from oc_meta.lib.master_of_regex import (
30 RE_COLON_AND_SPACES,
31 RE_MULTIPLE_SPACES,
32 RE_ONE_OR_MORE_SPACES,
33 RE_SEMICOLON_IN_PEOPLE_FIELD,
34 split_name_and_ids,
35)
36from oc_ocdm.counter_handler.counter_handler import CounterHandler
38if TYPE_CHECKING:
39 from rich.progress import Progress
42def _omid(meta: str) -> str:
43 return f"omid:{meta}"
46def _extract_ids_from_chunk(rows: list) -> Tuple[set, set, set]:
47 all_metavals = set()
48 all_identifiers = set()
49 all_vvis = set()
51 for row in rows:
52 metavals = set()
53 identifiers = set()
54 vvis = set()
55 venue_ids = set()
56 venue_metaid = None
58 if row["id"]:
59 id_list = RE_ONE_OR_MORE_SPACES.split(
60 RE_COLON_AND_SPACES.sub(":", row["id"])
61 )
62 idslist, metaval = Curator.clean_id_list(id_list, br=True)
63 if metaval:
64 metavals.add(_omid(metaval))
65 if idslist:
66 identifiers.update(idslist)
68 fields_with_an_id = []
69 for field in ["author", "editor", "publisher", "venue", "volume", "issue"]:
70 _, ids_str = split_name_and_ids(row[field])
71 if ids_str:
72 fields_with_an_id.append((field, ids_str.split()))
73 for field, field_ids in fields_with_an_id:
74 br = field in ["venue", "volume", "issue"]
75 field_idslist, field_metaval = Curator.clean_id_list(field_ids, br=br)
76 if field_metaval:
77 field_metaval = _omid(field_metaval)
78 else:
79 field_metaval = ""
80 if field_metaval:
81 metavals.add(field_metaval)
82 if field == "venue":
83 venue_metaid = field_metaval
84 if field_idslist:
85 venue_ids.update(field_idslist)
86 else:
87 if field_idslist:
88 identifiers.update(field_idslist)
90 if (venue_metaid or venue_ids) and (row["volume"] or row["issue"]):
91 vvi = (row["volume"], row["issue"], venue_metaid, tuple(sorted(venue_ids)))
92 vvis.add(vvi)
94 all_metavals.update(metavals)
95 all_identifiers.update(identifiers)
96 all_vvis.update(vvis)
98 return all_metavals, all_identifiers, all_vvis
101class Curator:
102 def __init__(
103 self,
104 data: List[dict],
105 ts: str,
106 prov_config: str,
107 counter_handler: CounterHandler,
108 base_iri: str = "https://w3id.org/oc/meta",
109 prefix: str = "060",
110 settings: dict | None = None,
111 silencer: list = [],
112 meta_config_path: str | None = None,
113 timer=None,
114 progress: Progress | None = None,
115 min_rows_parallel: int = 1000,
116 ):
117 self.timer = timer
118 self.progress = progress
119 self.settings = settings or {}
120 self.workers = self.settings.get("workers", 1)
121 self.finder = ResourceFinder(
122 ts,
123 base_iri,
124 settings=self.settings,
125 meta_config_path=meta_config_path,
126 workers=self.workers,
127 )
128 self.base_iri = base_iri
129 self.prov_config = prov_config
130 # Preliminary pass to clear volume and issue if id is present but venue is missing
131 for row in data:
132 if row["id"] and (row["volume"] or row["issue"]):
133 if not row["venue"]:
134 row["volume"] = ""
135 row["issue"] = ""
136 if not row["type"]:
137 row["type"] = "journal article"
138 self.data = [
139 {field: value.strip() for field, value in row.items()}
140 for row in data
141 if is_a_valid_row(row)
142 ]
143 self.prefix = prefix
144 self.counter_handler = counter_handler
146 self.entity_store = EntityStore()
147 self.ardict = {}
148 self.vvi = {}
149 self.remeta = dict()
150 self.wnb_cnt = 0
151 self.rowcnt = 0
152 self.preexisting_entities = set()
153 self.silencer = silencer
154 self.min_rows_parallel = min_rows_parallel
155 self.identifiers_only = self.settings.get("identifiers_only", False)
157 def _timed(self, name: str):
158 if self.timer:
159 return self.timer.timer(name)
160 return nullcontext()
162 def collect_identifiers(self):
163 return self._collect_identifiers_with_progress(task_id=None)
165 def _collect_identifiers_with_progress(self, task_id=None):
166 all_metavals = set()
167 all_idslist = set()
168 all_vvis = set()
170 total_rows = len(self.data)
171 if total_rows == 0:
172 return all_metavals, all_idslist, all_vvis
174 if total_rows > self.min_rows_parallel and self.workers > 1:
175 chunks = []
176 for i in range(0, total_rows, self.min_rows_parallel):
177 chunks.append(self.data[i : i + self.min_rows_parallel])
179 with ProcessPoolExecutor(
180 max_workers=self.workers,
181 mp_context=multiprocessing.get_context("forkserver"),
182 ) as executor:
183 for chunk_metavals, chunk_ids, chunk_vvis in executor.map(
184 _extract_ids_from_chunk, chunks
185 ):
186 all_metavals.update(chunk_metavals)
187 all_idslist.update(chunk_ids)
188 all_vvis.update(chunk_vvis)
189 if self.progress and task_id is not None:
190 self.progress.advance(
191 task_id, min(self.min_rows_parallel, total_rows)
192 )
193 else:
194 for row in self.data:
195 metavals, idslist, vvis = self.extract_identifiers_and_metavals(row)
196 all_metavals.update(metavals)
197 all_idslist.update(idslist)
198 all_vvis.update(vvis)
199 if self.progress and task_id is not None:
200 self.progress.advance(task_id)
202 return all_metavals, all_idslist, all_vvis
204 def extract_identifiers_and_metavals(self, row) -> Tuple[set, set, set]:
205 metavals = set()
206 identifiers = set()
207 vvis = set()
208 venue_ids = set()
209 venue_metaid = None
211 if row["id"]:
212 idslist, metaval = self.clean_id_list(
213 self.split_identifiers(row["id"]),
214 br=True,
215 )
216 id_metaval = _omid(metaval) if metaval else ""
217 if id_metaval:
218 metavals.add(id_metaval)
219 if idslist:
220 identifiers.update(idslist)
222 fields_with_an_id = []
223 for field in ["author", "editor", "publisher", "venue", "volume", "issue"]:
224 _, ids_str = split_name_and_ids(row[field])
225 if ids_str:
226 fields_with_an_id.append((field, ids_str.split()))
227 for field, field_ids in fields_with_an_id:
228 br = field in ["venue", "volume", "issue"]
229 field_idslist, field_metaval = self.clean_id_list(field_ids, br=br)
230 if field_metaval:
231 field_metaval = _omid(field_metaval)
232 else:
233 field_metaval = ""
234 if field_metaval:
235 metavals.add(field_metaval)
236 if field == "venue":
237 venue_metaid = field_metaval
238 if field_idslist:
239 venue_ids.update(field_idslist)
240 else:
241 if field_idslist:
242 identifiers.update(field_idslist)
244 if (venue_metaid or venue_ids) and (row["volume"] or row["issue"]):
245 vvi = (row["volume"], row["issue"], venue_metaid, tuple(sorted(venue_ids)))
246 vvis.add(vvi)
248 return metavals, identifiers, vvis
250 def split_identifiers(self, field_value):
251 return RE_ONE_OR_MORE_SPACES.split(RE_COLON_AND_SPACES.sub(":", field_value))
253 def curator(self, filename: str | None = None, path_csv: str | None = None):
254 total_rows = len(self.data)
256 # Phase 1: Collect identifiers and SPARQL prefetch
257 with self._timed("curation__collect_identifiers"):
258 task_collect = None
259 if self.progress:
260 task_collect = self.progress.add_task(
261 " [dim]Collecting identifiers[/dim]", total=total_rows
262 )
263 metavals, identifiers, vvis = self._collect_identifiers_with_progress(
264 task_id=task_collect,
265 )
266 if self.progress and task_collect is not None:
267 self.progress.remove_task(task_collect)
268 self.finder.get_everything_about_res(
269 metavals=metavals,
270 identifiers=identifiers,
271 vvis=vvis,
272 max_depth=1 if self.identifiers_only else 10,
273 progress=self.progress,
274 )
276 # Phase 2: Clean ID (loop over all rows)
277 with self._timed("curation__clean_id"):
278 task_clean_id = None
279 if self.progress:
280 task_clean_id = self.progress.add_task(
281 " [dim]Cleaning IDs[/dim]", total=total_rows
282 )
283 for row in self.data:
284 self.clean_id(row)
285 self.rowcnt += 1
286 if self.progress and task_clean_id is not None:
287 self.progress.advance(task_clean_id)
288 if self.progress and task_clean_id is not None:
289 self.progress.remove_task(task_clean_id)
291 # Phase 3: Merge duplicates + clean VVI/RA
292 with self._timed("curation__clean_vvi_ra"):
293 task_merge = None
294 if self.progress:
295 task_merge = self.progress.add_task(
296 " [dim]Merging duplicates[/dim]", total=total_rows
297 )
298 self.merge_duplicate_entities(task_id=task_merge)
299 if self.progress and task_merge is not None:
300 self.progress.remove_task(task_merge)
301 self.clean_metadata_without_id()
303 if not self.identifiers_only:
304 self.rowcnt = 0
305 task_vvi_ra = None
306 if self.progress:
307 task_vvi_ra = self.progress.add_task(
308 " [dim]Cleaning VVI and RA[/dim]", total=total_rows
309 )
310 for row in self.data:
311 self.clean_vvi(row)
312 self.clean_ra(row, "author")
313 self.clean_ra(row, "publisher")
314 self.clean_ra(row, "editor")
315 self.rowcnt += 1
316 if self.progress and task_vvi_ra is not None:
317 self.progress.advance(task_vvi_ra)
318 if self.progress and task_vvi_ra is not None:
319 self.progress.remove_task(task_vvi_ra)
321 # Phase 4: Metamaker (preexisting + meta_maker + enrich + dedupe)
322 with self._timed("curation__metamaker"):
323 task_metamaker = None
324 if self.progress:
325 task_metamaker = self.progress.add_task(
326 " [dim]Metamaker[/dim]", total=total_rows
327 )
328 self.get_preexisting_entities()
329 self.meta_maker()
330 self.enrich(task_id=task_metamaker)
331 if self.progress and task_metamaker is not None:
332 self.progress.remove_task(task_metamaker)
333 self.data = list({v["id"]: v for v in self.data}.values())
335 # Phase 5: CSV output (curated CSV indexer)
336 with self._timed("curation__csv_out"):
337 self.filename = filename
338 self.indexer(path_csv=path_csv)
340 def clean_id(self, row: Dict[str, str]) -> None:
341 """
342 The 'clean id()' function is executed for each CSV row.
343 In this process, any duplicates are detected by the IDs in the 'id' column.
344 For each line, a wannabeID or, if the bibliographic resource was found in the triplestore,
345 a MetaID is assigned.
346 Finally, this method enrich and clean the fields related to the
347 title, venue, volume, issue, page, publication date and type.
349 :params row: a dictionary representing a CSV row
350 :type row: Dict[str, str]
351 :returns: None -- This method modifies the input CSV row without returning it.
352 """
353 if row["title"]:
354 name = clean_title(
355 row["title"], bool(self.settings.get("normalize_titles", False))
356 )
357 else:
358 name = ""
359 metaval_ids_list = []
360 idslist: list = []
361 metaval = ""
362 if row["id"]:
363 idslist = RE_ONE_OR_MORE_SPACES.split(
364 RE_COLON_AND_SPACES.sub(":", row["id"])
365 )
366 idslist, metaval = self.clean_id_list(idslist, br=True)
367 id_metaval = _omid(metaval) if metaval else ""
368 metaval_ids_list.append((id_metaval, idslist))
369 fields_with_an_id = []
370 for field in ["author", "editor", "publisher", "venue", "volume", "issue"]:
371 _, ids_str = split_name_and_ids(row[field])
372 if ids_str:
373 fields_with_an_id.append((field, ids_str.split()))
374 for field, field_ids in fields_with_an_id:
375 br = field in ["venue", "volume", "issue"]
376 field_idslist, field_metaval = self.clean_id_list(field_ids, br=br)
377 if field_metaval:
378 field_metaval = _omid(field_metaval)
379 else:
380 field_metaval = ""
381 metaval_ids_list.append((field_metaval, field_idslist))
382 if row["id"]:
383 metaval = self.id_worker(
384 "id",
385 name,
386 idslist,
387 metaval,
388 ra_ent=False,
389 br_ent=True,
390 vvi_ent=False,
391 publ_entity=False,
392 )
393 else:
394 metaval = self.new_entity(name, "br")
395 row["title"] = self.entity_store.get_title(metaval)
396 row["id"] = metaval
398 def clean_metadata_without_id(self):
399 for row in self.data:
400 if row["page"]:
401 row["page"] = normalize_hyphens(row["page"])
402 if pub_date := row["pub_date"]:
403 row["pub_date"] = clean_date(normalize_hyphens(pub_date))
404 if row["type"]:
405 entity_type = RE_MULTIPLE_SPACES.sub(" ", row["type"].lower()).strip()
406 if entity_type == "edited book" or entity_type == "monograph":
407 entity_type = "book"
408 elif (
409 entity_type == "report series"
410 or entity_type == "standard series"
411 or entity_type == "proceedings series"
412 ):
413 entity_type = "series"
414 elif entity_type == "posted content":
415 entity_type = "web content"
416 if entity_type in VALID_ENTITY_TYPES:
417 row["type"] = entity_type
418 else:
419 row["type"] = ""
421 def clean_vvi(self, row: Dict[str, str]) -> None:
422 """
423 This method performs the deduplication process for venues, volumes and issues.
424 The acquired information is stored in the 'vvi' dictionary, that has the following format: ::
426 {
427 VENUE_IDENTIFIER: {
428 'issue': {SEQUENCE_IDENTIFIER: {'id': META_ID}},
429 'volume': {
430 SEQUENCE_IDENTIFIER: {
431 'id': META_ID,
432 'issue' {SEQUENCE_IDENTIFIER: {'id': META_ID}}
433 }
434 }
435 }
436 }
438 {
439 '4416': {
440 'issue': {},
441 'volume': {
442 '166': {'id': '4388', 'issue': {'4': {'id': '4389'}}},
443 '172': {'id': '4434',
444 'issue': {
445 '22': {'id': '4435'},
446 '20': {'id': '4436'},
447 '21': {'id': '4437'},
448 '19': {'id': '4438'}
449 }
450 }
451 }
452 }
453 }
455 :params row: a dictionary representing a CSV row
456 :type row: Dict[str, str]
457 :returns: None -- This method modifies the input CSV row without returning it.
458 """
459 if row["type"] not in {
460 "journal article",
461 "journal volume",
462 "journal issue",
463 } and (row["volume"] or row["issue"]):
464 row["volume"] = ""
465 row["issue"] = ""
466 clean_volume_and_issue(row=row)
467 vol_meta = None
468 br_type = row["type"]
469 volume = row["volume"]
470 issue = row["issue"]
471 br_id = row["id"]
472 venue = row["venue"]
473 # Venue
474 if venue:
475 # The data must be invalidated, because the resource is journal but a volume or an issue have also been specified
476 if br_type == "journal" and (volume or issue):
477 row["venue"] = ""
478 row["volume"] = ""
479 row["issue"] = ""
480 venue_name, venue_ids_str = split_name_and_ids(venue)
481 if venue_ids_str:
482 name = clean_title(
483 venue_name, bool(self.settings.get("normalize_titles", False))
484 )
485 idslist = RE_ONE_OR_MORE_SPACES.split(
486 RE_COLON_AND_SPACES.sub(":", venue_ids_str)
487 )
488 idslist, metaval = self.clean_id_list(idslist, br=True)
490 metaval = self.id_worker(
491 "venue",
492 name,
493 idslist,
494 metaval,
495 ra_ent=False,
496 br_ent=True,
497 vvi_ent=True,
498 publ_entity=False,
499 )
500 if metaval not in self.vvi:
501 ts_vvi = None
502 if "wannabe" not in metaval:
503 ts_vvi = self.finder.retrieve_venue_from_local_graph(metaval)
504 if "wannabe" in metaval or not ts_vvi:
505 self.vvi[metaval] = dict()
506 self.vvi[metaval]["volume"] = dict()
507 self.vvi[metaval]["issue"] = dict()
508 elif ts_vvi:
509 self.vvi[metaval] = ts_vvi
510 else:
511 name = clean_title(
512 venue_name or venue,
513 bool(self.settings.get("normalize_titles", False)),
514 )
515 metaval = self.new_entity(name, "br")
516 self.vvi[metaval] = dict()
517 self.vvi[metaval]["volume"] = dict()
518 self.vvi[metaval]["issue"] = dict()
519 row["venue"] = metaval
521 # Volume
522 if volume and (br_type == "journal issue" or br_type == "journal article"):
523 if volume in self.vvi[metaval]["volume"]:
524 vol_meta = self.vvi[metaval]["volume"][volume]["id"]
525 else:
526 vol_meta = self.new_entity("", "br")
527 self.vvi[metaval]["volume"][volume] = dict()
528 self.vvi[metaval]["volume"][volume]["id"] = vol_meta
529 self.vvi[metaval]["volume"][volume]["issue"] = dict()
530 elif volume and br_type == "journal volume":
531 # The data must be invalidated, because the resource is a journal volume but an issue has also been specified
532 if issue:
533 row["volume"] = ""
534 row["issue"] = ""
535 else:
536 vol_meta = br_id
537 self.volume_issue(
538 vol_meta, self.vvi[metaval]["volume"], volume, row
539 )
541 # Issue
542 if issue and br_type == "journal article":
543 row["issue"] = issue
544 if vol_meta:
545 if issue not in self.vvi[metaval]["volume"][volume]["issue"]:
546 issue_meta = self.new_entity("", "br")
547 self.vvi[metaval]["volume"][volume]["issue"][issue] = dict()
548 self.vvi[metaval]["volume"][volume]["issue"][issue]["id"] = (
549 issue_meta
550 )
551 else:
552 if issue not in self.vvi[metaval]["issue"]:
553 issue_meta = self.new_entity("", "br")
554 self.vvi[metaval]["issue"][issue] = dict()
555 self.vvi[metaval]["issue"][issue]["id"] = issue_meta
556 elif issue and br_type == "journal issue":
557 issue_meta = br_id
558 if vol_meta:
559 self.volume_issue(
560 issue_meta,
561 self.vvi[metaval]["volume"][volume]["issue"],
562 issue,
563 row,
564 )
565 else:
566 self.volume_issue(
567 issue_meta, self.vvi[metaval]["issue"], issue, row
568 )
570 else:
571 row["venue"] = ""
572 row["volume"] = ""
573 row["issue"] = ""
575 def clean_ra(self, row, col_name):
576 """
577 This method performs the deduplication process for responsible agents (authors, publishers and editors).
579 :params row: a dictionary representing a CSV row
580 :type row: Dict[str, str]
581 :params col_name: the CSV column name. It can be 'author', 'publisher', or 'editor'
582 :type col_name: str
583 :returns: None -- This method modifies self.ardict, self.radict, and self.idra, and returns None.
584 """
586 def get_br_metaval_to_check(row, col_name):
587 if col_name == "editor":
588 return get_edited_br_metaid(row, row["id"], row["venue"])
589 else:
590 return row["id"]
592 def get_br_metaval(br_metaval_to_check):
593 if (
594 br_metaval_to_check in self.entity_store
595 or br_metaval_to_check in self.vvi
596 ):
597 return br_metaval_to_check
598 return self.entity_store.find(br_metaval_to_check)
600 def initialize_ardict_entry(br_metaval):
601 if br_metaval not in self.ardict:
602 self.ardict[br_metaval] = {"author": [], "editor": [], "publisher": []}
604 def initialize_sequence(br_metaval, col_name):
605 sequence = []
606 if "wannabe" in br_metaval:
607 sequence = []
608 else:
609 sequence_found = self.finder.retrieve_ra_sequence_from_br_meta(
610 br_metaval, col_name
611 )
612 if sequence_found:
613 sequence = []
614 for agent in sequence_found:
615 for ar_metaid in agent:
616 ra_metaid = agent[ar_metaid][2]
617 sequence.append(tuple((ar_metaid, ra_metaid)))
618 if ra_metaid not in self.entity_store:
619 self.entity_store.add_entity(
620 ra_metaid, agent[ar_metaid][0]
621 )
622 for identifier in agent[ar_metaid][1]:
623 id_metaid = identifier[0]
624 literal = identifier[1]
625 if self.entity_store.get_id_metaid(literal) is None:
626 self.entity_store.set_id_metaid(literal, id_metaid)
627 if literal not in self.entity_store.get_ids(ra_metaid):
628 self.entity_store.add_id(ra_metaid, literal)
629 self.ardict[br_metaval][col_name].extend(sequence)
630 else:
631 sequence = []
632 return sequence
634 def parse_ra_list(row):
635 ra_list = RE_SEMICOLON_IN_PEOPLE_FIELD.split(row[col_name])
636 ra_list = clean_ra_list(ra_list)
637 return ra_list
639 def process_individual_ra(ra, sequence):
640 new_elem_seq = True
641 raw_name, ra_id = split_name_and_ids(ra)
642 name = clean_name(raw_name)
643 if not ra_id and sequence:
644 for _, ra_metaid in sequence:
645 if self.entity_store.get_title(ra_metaid) == name:
646 ra_id = "omid:" + str(ra_metaid)
647 new_elem_seq = False
648 break
649 return ra_id, name, new_elem_seq
651 if not row[col_name]:
652 return
654 br_metaval_to_check = get_br_metaval_to_check(row, col_name)
655 br_metaval = get_br_metaval(br_metaval_to_check)
656 initialize_ardict_entry(br_metaval)
658 sequence = self.ardict[br_metaval].get(col_name, [])
659 if not sequence:
660 sequence = initialize_sequence(br_metaval, col_name)
661 if col_name in self.silencer and sequence:
662 return
664 ra_list = parse_ra_list(row)
665 new_sequence = list()
667 for pos, ra in enumerate(ra_list):
668 ra_id, name, new_elem_seq = process_individual_ra(ra, sequence)
669 if ra_id:
670 ra_id_list = RE_ONE_OR_MORE_SPACES.split(
671 RE_COLON_AND_SPACES.sub(":", ra_id)
672 )
673 if sequence:
674 ar_ra = None
675 for el in sequence:
676 ra_metaid = el[1]
677 for literal in ra_id_list:
678 if literal in self.entity_store.get_ids(ra_metaid):
679 new_elem_seq = False
680 if "wannabe" not in ra_metaid:
681 ar_ra = ra_metaid
682 for pos, literal_value in enumerate(ra_id_list):
683 if "omid" in literal_value:
684 ra_id_list[pos] = ""
685 break
686 ra_id_list = list(filter(None, ra_id_list))
687 ra_id_list.append("omid:" + ar_ra)
688 if not ar_ra:
689 # new element
690 for ar_metaid, ra_metaid in sequence:
691 if self.entity_store.get_title(ra_metaid) == name:
692 new_elem_seq = False
693 if "wannabe" not in ra_metaid:
694 ar_ra = ra_metaid
695 for pos, i in enumerate(ra_id_list):
696 if "omid" in i:
697 ra_id_list[pos] = ""
698 break
699 ra_id_list = list(filter(None, ra_id_list))
700 ra_id_list.append("omid:" + ar_ra)
701 if col_name == "publisher":
702 ra_id_list, metaval = self.clean_id_list(ra_id_list, br=False)
703 metaval = self.id_worker(
704 "publisher",
705 name,
706 ra_id_list,
707 metaval,
708 ra_ent=True,
709 br_ent=False,
710 vvi_ent=False,
711 publ_entity=True,
712 )
713 else:
714 ra_id_list, metaval = self.clean_id_list(ra_id_list, br=False)
715 metaval = self.id_worker(
716 col_name,
717 name,
718 ra_id_list,
719 metaval,
720 ra_ent=True,
721 br_ent=False,
722 vvi_ent=False,
723 publ_entity=False,
724 )
725 if col_name != "publisher" and metaval in self.entity_store:
726 full_name: str = self.entity_store.get_title(metaval)
727 if "," in name and "," in full_name:
728 first_name = name.split(",")[1].strip()
729 if (
730 not full_name.split(",")[1].strip() and first_name
731 ): # first name found!
732 given_name = full_name.split(",")[0]
733 self.entity_store.set_title(
734 metaval, given_name + ", " + first_name
735 )
736 else:
737 metaval = self.new_entity(name, "ra")
738 if new_elem_seq:
739 role = "ar/" + self.prefix + str(self._add_number("ar"))
740 new_sequence.append(tuple((role, metaval)))
741 sequence.extend(new_sequence)
742 self.ardict[br_metaval][col_name] = sequence
744 @staticmethod
745 def clean_id_list(id_list: List[str], br: bool) -> Tuple[list, str]:
746 """
747 Clean IDs in the input list and check if there is a MetaID.
749 :params: id_list: a list of IDs
750 :type: id_list: List[str]
751 :params: br: True if the IDs in id_list refer to bibliographic resources, False otherwise
752 :type: br: bool
753 :returns: Tuple[list, str]: -- it returns a two-elements tuple, where the first element is the list of cleaned IDs, while the second is a MetaID (with prefix like "br/0601") if any was found.
754 """
755 metaid = ""
756 id_list = list(filter(None, id_list))
757 clean_set = set()
758 clean_list = []
760 for elem in id_list:
761 if elem in clean_set:
762 continue
763 clean_set.add(elem)
764 elem = normalize_hyphens(elem)
765 identifier = elem.split(":", 1)
766 schema = identifier[0].lower()
767 value = identifier[1]
769 if schema == "omid":
770 metaid = value
771 else:
772 normalized_id = normalize_id(elem)
773 if normalized_id:
774 clean_list.append(normalized_id)
776 meta_count = sum(1 for i in id_list if i.lower().startswith("omid"))
777 if meta_count > 1:
778 clean_list = [i for i in clean_list if not i.lower().startswith("omid")]
780 return clean_list, metaid
782 def conflict(
783 self, idslist: List[str], name: str, entity_type: str, col_name: str
784 ) -> str:
785 metaval = self.new_entity(name, entity_type)
786 for identifier in idslist:
787 self.entity_store.add_id(metaval, identifier)
788 if self.entity_store.get_id_metaid(identifier) is None:
789 schema_value = identifier.split(":", maxsplit=1)
790 found_metaid = self.finder.retrieve_metaid_from_id(
791 schema_value[0], schema_value[1]
792 )
793 if found_metaid:
794 self.entity_store.set_id_metaid(identifier, found_metaid)
795 else:
796 self.__update_id_count(identifier)
797 return metaval
799 def finder_sparql(self, list_to_find, br=True, ra=False, vvi=False, publ=False):
800 match_elem = list()
801 id_set = set()
802 res = None
803 for elem in list_to_find:
804 if len(match_elem) < 2:
805 identifier = elem.split(":", maxsplit=1)
806 value = identifier[1]
807 schema = identifier[0]
808 if br:
809 res = self.finder.retrieve_br_from_id(schema, value)
810 elif ra:
811 res = self.finder.retrieve_ra_from_id(schema, value)
812 if res:
813 for f in res:
814 if f[0] not in id_set:
815 match_elem.append(f)
816 id_set.add(f[0])
817 return match_elem
819 def ra_update(self, row: dict, br_key: str, col_name: str) -> None:
820 if row[col_name]:
821 sequence = self.ardict[br_key][col_name] if br_key in self.ardict else []
822 ras_list = list()
823 for _, ra_id in sequence:
824 ra_name = self.entity_store.get_title(ra_id)
825 ra_ids_with_omid = self.entity_store.get_ids(ra_id) | {_omid(ra_id)}
826 ra = self.build_name_ids_string(ra_name, ra_ids_with_omid)
827 ras_list.append(ra)
828 row[col_name] = "; ".join(ras_list)
830 @staticmethod
831 def build_name_ids_string(name: str, ids: set) -> str:
832 if name and ids:
833 return f"{name} [{' '.join(ids)}]"
834 elif name:
835 return name
836 elif ids:
837 return f"[{' '.join(ids)}]"
838 return ""
840 def _local_match(self, list_to_match: list, entity_type: str = "") -> dict:
841 def sort_key(entity_key: str) -> tuple:
842 if "wannabe_" in entity_key:
843 parts = entity_key.split("wannabe_")
844 return (0, int(parts[1]))
845 return (1, entity_key)
847 entity_prefix = f"{entity_type}/" if entity_type else ""
848 match_elem: dict[str, list] = {"existing": [], "wannabe": []}
849 seen: set[str] = set()
850 for identifier in list_to_match:
851 entities = self.entity_store.find_entities(identifier)
852 for entity_key in sorted(entities, key=sort_key):
853 if entity_prefix and not entity_key.startswith(entity_prefix):
854 continue
855 if entity_key not in seen:
856 seen.add(entity_key)
857 if "wannabe" in entity_key:
858 match_elem["wannabe"].append(entity_key)
859 else:
860 match_elem["existing"].append(entity_key)
861 return match_elem
863 def __tree_traverse(self, tree: dict, key: str, values: List[Tuple]) -> None:
864 for k, v in tree.items():
865 if k == key:
866 values.append(v)
867 elif isinstance(v, dict):
868 found = self.__tree_traverse(v, key, values)
869 if found is not None:
870 values.append(found)
872 def get_preexisting_entities(self) -> None:
873 for entity_metaid in self.entity_store:
874 if "wannabe" not in entity_metaid:
875 self.preexisting_entities.add(entity_metaid)
876 for entity_id_literal in self.entity_store.get_ids(entity_metaid):
877 preexisting_entity_id_metaid = self.entity_store.get_id_metaid(
878 entity_id_literal
879 )
880 if preexisting_entity_id_metaid:
881 self.preexisting_entities.add(preexisting_entity_id_metaid)
882 for _, roles in self.ardict.items():
883 for _, ar_ras in roles.items():
884 for ar_ra in ar_ras:
885 if "wannabe" not in ar_ra[1]:
886 self.preexisting_entities.add(ar_ra[0])
887 for venue_metaid, vi in self.vvi.items():
888 if "wannabe" not in venue_metaid:
889 wannabe_preexisting_vis = list()
890 self.__tree_traverse(vi, "id", wannabe_preexisting_vis)
891 self.preexisting_entities.update(
892 {
893 vi_metaid
894 for vi_metaid in wannabe_preexisting_vis
895 if "wannabe" not in vi_metaid
896 }
897 )
898 for _, re_metaid in self.remeta.items():
899 re_id = re_metaid[0]
900 if not re_id.startswith("re/"):
901 re_id = f"re/{re_id}"
902 self.preexisting_entities.add(re_id)
904 def meta_maker(self):
905 """
906 Converts temporary wannabe identifiers to final MetaIDs.
907 Assigns final MetaIDs to wannabe entities and resolves ardict.
908 VolIss is also resolved from vvi.
909 """
910 for identifier in list(self.entity_store):
911 if identifier.startswith("br/") and "wannabe" in identifier:
912 count = self._add_number("br")
913 target_meta = f"br/{self.prefix}{count}"
914 self.entity_store.assign_meta(identifier, target_meta)
915 elif identifier.startswith("ra/") and "wannabe" in identifier:
916 count = self._add_number("ra")
917 target_meta = f"ra/{self.prefix}{count}"
918 self.entity_store.assign_meta(identifier, target_meta)
920 resolved_ardict: dict[str, dict[str, list]] = {}
921 for ar_id in self.ardict:
922 br_key = self.entity_store.find(ar_id)
923 if br_key not in resolved_ardict:
924 resolved_ardict[br_key] = {"author": [], "editor": [], "publisher": []}
925 for role_type in ["author", "editor", "publisher"]:
926 for ar_metaid, agent_id in self.ardict[ar_id][role_type]:
927 resolved_ra_metaid = self.entity_store.find(agent_id)
928 resolved_ardict[br_key][role_type].append(
929 (ar_metaid, resolved_ra_metaid)
930 )
931 self.ardict = resolved_ardict
933 self.VolIss = dict()
934 if self.vvi:
935 for venue_meta in self.vvi:
936 venue_issue = self.vvi[venue_meta]["issue"]
937 if venue_issue:
938 for issue in venue_issue:
939 issue_id = venue_issue[issue]["id"]
940 if "wannabe" in issue_id:
941 self.vvi[venue_meta]["issue"][issue]["id"] = str(
942 self.entity_store.find(issue_id)
943 )
945 venue_volume = self.vvi[venue_meta]["volume"]
946 if venue_volume:
947 for volume in venue_volume:
948 volume_id = venue_volume[volume]["id"]
949 if "wannabe" in volume_id:
950 self.vvi[venue_meta]["volume"][volume]["id"] = str(
951 self.entity_store.find(volume_id)
952 )
953 if venue_volume[volume]["issue"]:
954 volume_issue = venue_volume[volume]["issue"]
955 for issue in volume_issue:
956 volume_issue_id = volume_issue[issue]["id"]
957 if "wannabe" in volume_issue_id:
958 self.vvi[venue_meta]["volume"][volume]["issue"][
959 issue
960 ]["id"] = str(
961 self.entity_store.find(volume_issue_id)
962 )
963 if "wannabe" in venue_meta:
964 br_meta = self.entity_store.find(venue_meta)
965 self._merge_VolIss_with_vvi(br_meta, venue_meta)
966 else:
967 self._merge_VolIss_with_vvi(venue_meta, venue_meta)
969 def enrich(self, task_id=None):
970 """
971 This method replaces the wannabeID placeholders with the
972 actual data and MetaIDs as a result of the deduplication process.
973 """
974 for row in self.data:
975 metaid = row["id"]
976 if "wannabe" in row["id"]:
977 metaid = self.entity_store.find(row["id"])
978 if row["page"] and (metaid not in self.remeta):
979 re_meta = self.finder.retrieve_re_from_br_meta(metaid)
980 if re_meta:
981 self.remeta[metaid] = re_meta
982 row["page"] = re_meta[1]
983 else:
984 count = "re/" + self.prefix + str(self._add_number("re"))
985 page = row["page"]
986 self.remeta[metaid] = (count, page)
987 row["page"] = page
988 elif metaid in self.remeta:
989 row["page"] = self.remeta[metaid][1]
990 ids_with_omid = self.entity_store.get_ids(metaid) | {_omid(metaid)}
991 row["id"] = " ".join(ids_with_omid)
992 row["title"] = self.entity_store.get_title(metaid)
993 venue_metaid = None
994 if row["venue"]:
995 venue = row["venue"]
996 if "wannabe" in venue:
997 venue_metaid = self.entity_store.find(venue)
998 else:
999 venue_metaid = venue
1000 venue_ids_with_omid = self.entity_store.get_ids(venue_metaid) | {
1001 _omid(venue_metaid)
1002 }
1003 row["venue"] = self.build_name_ids_string(
1004 self.entity_store.get_title(venue_metaid), venue_ids_with_omid
1005 )
1006 br_key_for_editor = get_edited_br_metaid(row, metaid, venue_metaid)
1007 self.ra_update(row, metaid, "author")
1008 self.ra_update(row, metaid, "publisher")
1009 self.ra_update(row, br_key_for_editor, "editor")
1010 if self.progress and task_id is not None:
1011 self.progress.advance(task_id)
1013 @staticmethod
1014 def name_check(ts_name, name):
1015 if "," in ts_name:
1016 names = ts_name.split(",")
1017 if names[0] and not names[1].strip():
1018 if "," in name:
1019 gname = name.split(", ")[1]
1020 if gname.strip():
1021 ts_name = names[0] + ", " + gname
1022 return ts_name
1024 def _read_number(self, entity_type: str) -> int:
1025 return self.counter_handler.read_counter(
1026 entity_type, supplier_prefix=self.prefix
1027 )
1029 def _add_number(self, entity_type: str) -> int:
1030 return self.counter_handler.increment_counter(
1031 entity_type, supplier_prefix=self.prefix
1032 )
1034 def __update_id_and_entity_store(
1035 self,
1036 existing_ids: list,
1037 metaval: str,
1038 ) -> None:
1039 for identifier in existing_ids:
1040 if self.entity_store.get_id_metaid(identifier[1]) is None:
1041 self.entity_store.set_id_metaid(identifier[1], identifier[0])
1042 if identifier[1] not in self.entity_store.get_ids(metaval):
1043 self.entity_store.add_id(metaval, identifier[1])
1045 def indexer(self, path_csv: str | None = None) -> None:
1046 """
1047 Transform internal dicts (idra, idbr, ardict, remeta) to list-of-dicts format
1048 for Creator consumption. Optionally saves the enriched CSV file.
1050 :params path_csv: Directory path for the enriched CSV output (optional)
1051 :type path_csv: str
1052 """
1053 self.index_id_ra = list()
1054 self.index_id_br = list()
1055 id_metaids = self.entity_store.get_id_metaids()
1056 for literal, metaid in id_metaids.items():
1057 entities = self.entity_store.find_entities(literal)
1058 has_br = any(e.startswith("br/") for e in entities)
1059 has_ra = any(e.startswith("ra/") for e in entities)
1060 if has_br:
1061 self.index_id_br.append({"id": str(literal), "meta": str(metaid)})
1062 if has_ra:
1063 self.index_id_ra.append({"id": str(literal), "meta": str(metaid)})
1064 if not self.index_id_br:
1065 self.index_id_br.append({"id": "", "meta": ""})
1066 if not self.index_id_ra:
1067 self.index_id_ra.append({"id": "", "meta": ""})
1068 self.ar_index = list()
1069 if self.ardict:
1070 for metaid in self.ardict:
1071 index = dict()
1072 index["meta"] = metaid
1073 for role in self.ardict[metaid]:
1074 list_ar = list()
1075 for ar, ra in self.ardict[metaid][role]:
1076 list_ar.append(str(ar) + ", " + str(ra))
1077 index[role] = "; ".join(list_ar)
1078 self.ar_index.append(index)
1079 else:
1080 row = dict()
1081 row["meta"] = ""
1082 row["author"] = ""
1083 row["editor"] = ""
1084 row["publisher"] = ""
1085 self.ar_index.append(row)
1086 self.re_index = list()
1087 if self.remeta:
1088 for x in self.remeta:
1089 r = dict()
1090 r["br"] = x
1091 r["re"] = str(self.remeta[x][0])
1092 self.re_index.append(r)
1093 else:
1094 row = dict()
1095 row["br"] = ""
1096 row["re"] = ""
1097 self.re_index.append(row)
1098 if self.filename and path_csv and self.data:
1099 name = self.filename + ".csv"
1100 data_file = os.path.join(path_csv, name)
1101 write_csv(data_file, self.data)
1103 def _merge_VolIss_with_vvi(
1104 self, VolIss_venue_meta: str, vvi_venue_meta: str
1105 ) -> None:
1106 if VolIss_venue_meta in self.VolIss:
1107 for vvi_v in self.vvi[vvi_venue_meta]["volume"]:
1108 if vvi_v in self.VolIss[VolIss_venue_meta]["volume"]:
1109 self.VolIss[VolIss_venue_meta]["volume"][vvi_v]["issue"].update(
1110 self.vvi[vvi_venue_meta]["volume"][vvi_v]["issue"]
1111 )
1112 else:
1113 self.VolIss[VolIss_venue_meta]["volume"][vvi_v] = self.vvi[
1114 vvi_venue_meta
1115 ]["volume"][vvi_v]
1116 self.VolIss[VolIss_venue_meta]["issue"].update(
1117 self.vvi[vvi_venue_meta]["issue"]
1118 )
1119 else:
1120 self.VolIss[VolIss_venue_meta] = self.vvi[vvi_venue_meta]
1122 def __update_id_count(self, identifier):
1123 schema, value = identifier.split(":", maxsplit=1)
1124 existing_metaid = self.finder.retrieve_metaid_from_id(schema, value)
1126 if existing_metaid:
1127 self.entity_store.set_id_metaid(identifier, existing_metaid)
1128 else:
1129 count = self._add_number("id")
1130 self.entity_store.set_id_metaid(identifier, f"id/{self.prefix}{count}")
1132 def merge(
1133 self,
1134 metaval: str,
1135 old_meta: str,
1136 temporary_name: str,
1137 ) -> None:
1138 for sid in self.entity_store.get_ids(old_meta):
1139 self.entity_store.add_id(metaval, sid)
1140 if not self.entity_store.get_title(metaval):
1141 title = self.entity_store.get_title(old_meta) or temporary_name
1142 self.entity_store.set_title(metaval, title)
1143 self.entity_store.update_id_entity(old_meta, metaval)
1144 self.entity_store.merge(metaval, old_meta)
1145 self.entity_store.remove_entity(old_meta)
1147 def merge_entities_in_csv(
1148 self,
1149 idslist: list,
1150 metaval: str,
1151 name: str,
1152 ) -> None:
1153 entity_type = metaval.split("/")[0] if "/" in metaval else ""
1154 found_others = self._local_match(idslist, entity_type)
1155 if found_others["wannabe"]:
1156 for old_meta in found_others["wannabe"]:
1157 self.merge(metaval, old_meta, name)
1158 entry_ids = self.entity_store.get_ids(metaval)
1159 for identifier in idslist:
1160 if identifier not in entry_ids:
1161 self.entity_store.add_id(metaval, identifier)
1162 if self.entity_store.get_id_metaid(identifier) is None:
1163 self.__update_id_count(identifier)
1164 if not self.entity_store.get_title(metaval) and name:
1165 self.entity_store.set_title(metaval, name)
1167 def id_worker(
1168 self,
1169 col_name,
1170 name,
1171 idslist: List[str],
1172 metaval: str,
1173 ra_ent=False,
1174 br_ent=False,
1175 vvi_ent=False,
1176 publ_entity=False,
1177 ):
1178 entity_type = "ra" if ra_ent else "br"
1179 if metaval:
1180 if metaval in self.entity_store:
1181 self.merge_entities_in_csv(idslist, metaval, name)
1182 else:
1183 found_meta_ts: tuple[str, list[tuple[str, str]], bool] = ("", [], False)
1184 if ra_ent:
1185 found_meta_ts = self.finder.retrieve_ra_from_meta(metaval)
1186 elif br_ent:
1187 found_meta_ts = self.finder.retrieve_br_from_meta(metaval)
1188 if found_meta_ts[2]:
1189 title = (
1190 self.name_check(found_meta_ts[0], name)
1191 if col_name in ("author", "editor")
1192 else found_meta_ts[0]
1193 )
1194 self.entity_store.add_entity(metaval, title)
1195 existing_ids = found_meta_ts[1]
1196 self.__update_id_and_entity_store(existing_ids, metaval)
1197 self.merge_entities_in_csv(idslist, metaval, name)
1198 else:
1199 metaid_uri = f"{self.base_iri}/{metaval}"
1200 merged_metaval = self.finder.retrieve_metaid_from_merged_entity(
1201 metaid_uri=metaid_uri, prov_config=self.prov_config
1202 )
1203 metaval = (
1204 f"{entity_type}/{merged_metaval}" if merged_metaval else ""
1205 )
1206 if idslist and not metaval:
1207 local_match = self._local_match(idslist, entity_type)
1208 if local_match["existing"]:
1209 if len(local_match["existing"]) > 1:
1210 return self.conflict(idslist, name, entity_type, col_name)
1211 elif len(local_match["existing"]) == 1:
1212 metaval = str(local_match["existing"][0])
1213 entry_ids = self.entity_store.get_ids(metaval)
1214 suspect_ids = [i for i in idslist if i not in entry_ids]
1215 if suspect_ids:
1216 sparql_match = self.finder_sparql(
1217 suspect_ids,
1218 br=br_ent,
1219 ra=ra_ent,
1220 vvi=vvi_ent,
1221 publ=publ_entity,
1222 )
1223 if len(sparql_match) > 1:
1224 return self.conflict(idslist, name, entity_type, col_name)
1225 elif local_match["wannabe"]:
1226 metaval = str(local_match["wannabe"].pop(0))
1227 for old_meta in local_match["wannabe"]:
1228 self.merge(metaval, old_meta, name)
1229 entry_ids = self.entity_store.get_ids(metaval)
1230 suspect_ids = [i for i in idslist if i not in entry_ids]
1231 if suspect_ids:
1232 sparql_match = self.finder_sparql(
1233 suspect_ids, br=br_ent, ra=ra_ent, vvi=vvi_ent, publ=publ_entity
1234 )
1235 if sparql_match:
1236 # if 'wannabe' not in metaval or len(sparql_match) > 1:
1237 # # Two entities previously disconnected on the triplestore now become connected
1238 # # !
1239 # return self.conflict(idslist, name, id_dict, col_name)
1240 # else:
1241 # Collect all existing IDs from all matches
1242 existing_ids = []
1243 for match in sparql_match:
1244 existing_ids.extend(match[2])
1246 # new_idslist = [x[1] for x in existing_ids]
1247 # new_sparql_match = self.finder_sparql(new_idslist, br=br_ent, ra=ra_ent, vvi=vvi_ent, publ=publ_entity)
1248 # if len(new_sparql_match) > 1:
1249 # # Two entities previously disconnected on the triplestore now become connected
1250 # # !
1251 # return self.conflict(idslist, name, id_dict, col_name)
1252 # else:
1253 # 4 Merge data from EntityA (CSV) with data from EntityX (CSV) (it has already happened in # 5), update both with data from EntityA (RDF)
1254 old_metaval = metaval
1255 metaval = sparql_match[0][0]
1256 self.entity_store.add_entity(metaval, sparql_match[0][1] or "")
1257 self.__update_id_and_entity_store(existing_ids, metaval)
1258 self.merge(metaval, old_metaval, sparql_match[0][1])
1259 else:
1260 sparql_match = self.finder_sparql(
1261 idslist, br=br_ent, ra=ra_ent, vvi=vvi_ent, publ=publ_entity
1262 )
1263 # if len(sparql_match) > 1:
1264 # # !
1265 # return self.conflict(idslist, name, id_dict, col_name)
1266 # elif len(sparql_match) == 1:
1267 if sparql_match:
1268 # Collect all existing IDs from all matches
1269 existing_ids = []
1270 for match in sparql_match:
1271 existing_ids.extend(match[2])
1273 # new_idslist = [x[1] for x in existing_ids]
1274 # new_sparql_match = self.finder_sparql(new_idslist, br=br_ent, ra=ra_ent, vvi=vvi_ent, publ=publ_entity)
1275 # if len(new_sparql_match) > 1:
1276 # # Two entities previously disconnected on the triplestore now become connected
1277 # # !
1278 # return self.conflict(idslist, name, id_dict, col_name)
1279 # 2 Retrieve EntityA data in triplestore to update EntityA inside CSV
1280 # 3 CONFLICT beteen MetaIDs. MetaID specified in EntityA inside CSV has precedence.
1281 # elif len(new_sparql_match) == 1:
1282 metaval = sparql_match[0][0]
1283 title = (
1284 self.name_check(sparql_match[0][1], name)
1285 if col_name in ("author", "editor")
1286 else sparql_match[0][1]
1287 )
1288 self.entity_store.add_entity(metaval, title or name)
1289 self.__update_id_and_entity_store(existing_ids, metaval)
1290 else:
1291 # 1 EntityA is a new one
1292 metaval = self.new_entity(name, entity_type)
1293 entry_ids = self.entity_store.get_ids(metaval)
1294 for identifier in idslist:
1295 if self.entity_store.get_id_metaid(identifier) is None:
1296 self.__update_id_count(identifier)
1297 if identifier not in entry_ids:
1298 self.entity_store.add_id(metaval, identifier)
1299 if not self.entity_store.get_title(metaval) and name:
1300 self.entity_store.set_title(metaval, name)
1301 # 1 EntityA is a new one
1302 if not idslist and not metaval:
1303 metaval = self.new_entity(name, entity_type)
1304 return metaval
1306 def new_entity(self, name: str, entity_type: str) -> str:
1307 metaval = f"{entity_type}/wannabe_{self.wnb_cnt}"
1308 self.wnb_cnt += 1
1309 self.entity_store.add_entity(metaval, name)
1310 return metaval
1312 def volume_issue(
1313 self,
1314 meta: str,
1315 path: dict,
1316 value: str,
1317 row: Dict[str, str],
1318 ) -> None:
1319 if "wannabe" not in meta:
1320 if value in path:
1321 if "wannabe" in path[value]["id"]:
1322 old_meta = path[value]["id"]
1323 self.merge(meta, old_meta, row["title"])
1324 path[value]["id"] = meta
1325 else:
1326 path[value] = {"id": meta}
1327 if "issue" not in path:
1328 path[value]["issue"] = {}
1329 else:
1330 if value in path:
1331 if "wannabe" in path[value]["id"]:
1332 old_meta = path[value]["id"]
1333 if meta != old_meta:
1334 self.merge(meta, old_meta, row["title"])
1335 path[value]["id"] = meta
1336 else:
1337 old_meta = path[value]["id"]
1338 if "wannabe" not in old_meta and old_meta not in self.entity_store:
1339 br4dict = self.finder.retrieve_br_from_meta(old_meta)
1340 self.entity_store.add_entity(
1341 old_meta, br4dict[0] if br4dict else ""
1342 )
1343 if br4dict:
1344 for x in br4dict[1]:
1345 identifier = x[1]
1346 self.entity_store.add_id(old_meta, identifier)
1347 if self.entity_store.get_id_metaid(identifier) is None:
1348 self.entity_store.set_id_metaid(identifier, x[0])
1349 self.merge(old_meta, meta, row["title"])
1350 else:
1351 path[value] = {"id": meta}
1352 if "issue" not in path:
1353 path[value]["issue"] = {}
1355 def merge_duplicate_entities(self, task_id=None) -> None:
1356 """
1357 Merges duplicate entities and propagates data from existing entities to related rows.
1358 For rows referencing existing entities, retrieves data from triplestore (via equalizer)
1359 and updates all related rows (those with matching IDs or merged entity references).
1360 """
1361 # Build index mapping row IDs to row indices for O(1) lookup
1362 id_to_indices: dict[str, list[int]] = {}
1363 for idx, row in enumerate(self.data):
1364 row_id = row["id"]
1365 if row_id not in id_to_indices:
1366 id_to_indices[row_id] = []
1367 id_to_indices[row_id].append(idx)
1369 self.rowcnt = 0
1370 for row in self.data:
1371 row_id = row["id"]
1372 if "wannabe" not in row_id:
1373 if not self.identifiers_only:
1374 self.equalizer(row, row_id)
1375 related_indices: set[int] = set()
1376 if row_id in id_to_indices:
1377 related_indices.update(id_to_indices[row_id])
1378 for other_id in self.entity_store.get_merged(row_id):
1379 if other_id in id_to_indices:
1380 related_indices.update(id_to_indices[other_id])
1381 related_indices.discard(self.rowcnt)
1382 for other_idx in related_indices:
1383 other_row = self.data[other_idx]
1384 for field in row:
1385 if row[field] and row[field] != other_row[field]:
1386 other_row[field] = row[field]
1387 if self.progress and task_id is not None:
1388 self.progress.advance(task_id)
1389 self.rowcnt += 1
1391 def extract_name_and_ids(self, venue_str: str) -> Tuple[str, List[str]]:
1392 """
1393 Extracts the name and IDs from the venue string.
1395 :params venue_str: the venue string
1396 :type venue_str: str
1397 :returns: Tuple[str, List[str]] -- the name and list of IDs extracted from the venue string
1398 """
1399 name, ids_str = split_name_and_ids(venue_str)
1400 return name.strip(), ids_str.split()
1402 def equalizer(self, row: Dict[str, str], metaval: str) -> None:
1403 """
1404 Given a CSV row and its MetaID, equates the information present in the CSV
1405 with that present on the triplestore. Triplestore data takes precedence.
1406 """
1407 known_data = self.finder.retrieve_br_info_from_meta(metaval)
1408 try:
1409 known_data["author"] = self.__get_resp_agents(metaval, "author")
1410 except ValueError:
1411 print(row)
1412 raise (ValueError)
1413 known_data["editor"] = self.__get_resp_agents(metaval, "editor")
1414 known_data["publisher"] = self.finder.retrieve_publisher_from_br_metaid(metaval)
1415 for datum in ["pub_date", "type", "volume", "issue"]:
1416 if known_data[datum]:
1417 row[datum] = known_data[datum]
1418 for datum in ["author", "editor", "publisher"]:
1419 if known_data[datum] and not row[datum]:
1420 row[datum] = known_data[datum]
1421 if known_data["venue"]:
1422 current_venue = row["venue"]
1423 known_venue = known_data["venue"]
1425 if current_venue:
1426 current_venue_name, current_venue_ids = self.extract_name_and_ids(
1427 current_venue
1428 )
1429 known_venue_name, known_venue_ids = self.extract_name_and_ids(
1430 known_venue
1431 )
1433 current_venue_ids_set = set(current_venue_ids)
1434 known_venue_ids_set = set(known_venue_ids)
1436 common_ids = current_venue_ids_set.intersection(known_venue_ids_set)
1438 if common_ids:
1439 merged_ids = current_venue_ids_set.union(known_venue_ids_set)
1440 row["venue"] = (
1441 f"{known_venue_name} [{' '.join(sorted(merged_ids))}]"
1442 )
1443 else:
1444 row["venue"] = known_venue
1445 else:
1446 row["venue"] = known_venue
1447 if known_data["page"]:
1448 row["page"] = known_data["page"][1]
1449 self.remeta[metaval] = known_data["page"]
1451 def __get_resp_agents(self, metaid: str, column: str) -> str:
1452 resp_agents = self.finder.retrieve_ra_sequence_from_br_meta(metaid, column)
1453 output = ""
1454 if resp_agents:
1455 full_resp_agents = list()
1456 for item in resp_agents:
1457 for _, resp_agent in item.items():
1458 author_name = resp_agent[0]
1459 ids = [_omid(resp_agent[2])]
1460 ids.extend([id[1] for id in resp_agent[1]])
1461 author_ids = "[" + " ".join(ids) + "]"
1462 full_resp_agent = author_name + " " + author_ids
1463 full_resp_agents.append(full_resp_agent)
1464 output = "; ".join(full_resp_agents)
1465 return output
1468def is_a_valid_row(row: Dict[str, str]) -> bool:
1469 """
1470 This method discards invalid rows in the input CSV file.
1472 :params row: a dictionary representing a CSV row
1473 :type row: Dict[str, str]
1474 :returns: bool -- This method returns True if the row is valid, False if it is invalid.
1475 """
1476 br_type = " ".join((row["type"].lower()).split())
1477 br_title = row["title"]
1478 br_volume = row["volume"]
1479 br_issue = row["issue"]
1480 br_venue = row["venue"]
1481 if row["id"]:
1482 if (br_volume or br_issue) and (not br_type or not br_venue):
1483 return False
1484 return True
1485 if all(not row[value] for value in row):
1486 return False
1487 br_author = row["author"]
1488 br_editor = row["editor"]
1489 br_pub_date = row["pub_date"]
1490 if not br_type or br_type in {
1491 "book",
1492 "data file",
1493 "dataset",
1494 "dissertation",
1495 "edited book",
1496 "journal article",
1497 "monograph",
1498 "other",
1499 "peer review",
1500 "posted content",
1501 "web content",
1502 "proceedings article",
1503 "report",
1504 "reference book",
1505 }:
1506 is_a_valid_row = (
1507 True if br_title and br_pub_date and (br_author or br_editor) else False
1508 )
1509 elif br_type in {
1510 "book chapter",
1511 "book part",
1512 "book section",
1513 "book track",
1514 "component",
1515 "reference entry",
1516 }:
1517 is_a_valid_row = True if br_title and br_venue else False
1518 elif br_type in {
1519 "book series",
1520 "book set",
1521 "journal",
1522 "proceedings",
1523 "proceedings series",
1524 "report series",
1525 "standard",
1526 "standard series",
1527 }:
1528 is_a_valid_row = True if br_title else False
1529 elif br_type == "journal volume":
1530 is_a_valid_row = True if br_venue and (br_volume or br_title) else False
1531 elif br_type == "journal issue":
1532 is_a_valid_row = True if br_venue and (br_issue or br_title) else False
1533 else:
1534 is_a_valid_row = False
1535 return is_a_valid_row
1538def get_edited_br_metaid(row: dict, metaid: str, venue_metaid: str | None) -> str:
1539 if (
1540 row["author"]
1541 and row["venue"]
1542 and row["type"] in CONTAINER_EDITOR_TYPES
1543 and venue_metaid
1544 ):
1545 return venue_metaid
1546 return metaid