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

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 

7 

8from __future__ import annotations 

9 

10import multiprocessing 

11import os 

12from concurrent.futures import ProcessPoolExecutor 

13from contextlib import nullcontext 

14from typing import TYPE_CHECKING, Dict, List, Tuple 

15 

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 

37 

38if TYPE_CHECKING: 

39 from rich.progress import Progress 

40 

41 

42def _omid(meta: str) -> str: 

43 return f"omid:{meta}" 

44 

45 

46def _extract_ids_from_chunk(rows: list) -> Tuple[set, set, set]: 

47 all_metavals = set() 

48 all_identifiers = set() 

49 all_vvis = set() 

50 

51 for row in rows: 

52 metavals = set() 

53 identifiers = set() 

54 vvis = set() 

55 venue_ids = set() 

56 venue_metaid = None 

57 

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) 

67 

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) 

89 

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) 

93 

94 all_metavals.update(metavals) 

95 all_identifiers.update(identifiers) 

96 all_vvis.update(vvis) 

97 

98 return all_metavals, all_identifiers, all_vvis 

99 

100 

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 

145 

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) 

156 

157 def _timed(self, name: str): 

158 if self.timer: 

159 return self.timer.timer(name) 

160 return nullcontext() 

161 

162 def collect_identifiers(self): 

163 return self._collect_identifiers_with_progress(task_id=None) 

164 

165 def _collect_identifiers_with_progress(self, task_id=None): 

166 all_metavals = set() 

167 all_idslist = set() 

168 all_vvis = set() 

169 

170 total_rows = len(self.data) 

171 if total_rows == 0: 

172 return all_metavals, all_idslist, all_vvis 

173 

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]) 

178 

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) 

201 

202 return all_metavals, all_idslist, all_vvis 

203 

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 

210 

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) 

221 

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) 

243 

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) 

247 

248 return metavals, identifiers, vvis 

249 

250 def split_identifiers(self, field_value): 

251 return RE_ONE_OR_MORE_SPACES.split(RE_COLON_AND_SPACES.sub(":", field_value)) 

252 

253 def curator(self, filename: str | None = None, path_csv: str | None = None): 

254 total_rows = len(self.data) 

255 

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 ) 

275 

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) 

290 

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() 

302 

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) 

320 

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()) 

334 

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) 

339 

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. 

348 

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 

397 

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"] = "" 

420 

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: :: 

425 

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 } 

437 

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 } 

454 

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) 

489 

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 

520 

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 ) 

540 

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 ) 

569 

570 else: 

571 row["venue"] = "" 

572 row["volume"] = "" 

573 row["issue"] = "" 

574 

575 def clean_ra(self, row, col_name): 

576 """ 

577 This method performs the deduplication process for responsible agents (authors, publishers and editors). 

578 

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 """ 

585 

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"] 

591 

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) 

599 

600 def initialize_ardict_entry(br_metaval): 

601 if br_metaval not in self.ardict: 

602 self.ardict[br_metaval] = {"author": [], "editor": [], "publisher": []} 

603 

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 

633 

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 

638 

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 

650 

651 if not row[col_name]: 

652 return 

653 

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) 

657 

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 

663 

664 ra_list = parse_ra_list(row) 

665 new_sequence = list() 

666 

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 

743 

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. 

748 

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 = [] 

759 

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] 

768 

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) 

775 

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")] 

779 

780 return clean_list, metaid 

781 

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 

798 

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 

818 

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) 

829 

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 "" 

839 

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) 

846 

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 

862 

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) 

871 

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) 

903 

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) 

919 

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 

932 

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 ) 

944 

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) 

968 

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) 

1012 

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 

1023 

1024 def _read_number(self, entity_type: str) -> int: 

1025 return self.counter_handler.read_counter( 

1026 entity_type, supplier_prefix=self.prefix 

1027 ) 

1028 

1029 def _add_number(self, entity_type: str) -> int: 

1030 return self.counter_handler.increment_counter( 

1031 entity_type, supplier_prefix=self.prefix 

1032 ) 

1033 

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]) 

1044 

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. 

1049 

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) 

1102 

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] 

1121 

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) 

1125 

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}") 

1131 

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) 

1146 

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) 

1166 

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]) 

1245 

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]) 

1272 

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 

1305 

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 

1311 

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"] = {} 

1354 

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) 

1368 

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 

1390 

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. 

1394 

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() 

1401 

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"] 

1424 

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 ) 

1432 

1433 current_venue_ids_set = set(current_venue_ids) 

1434 known_venue_ids_set = set(known_venue_ids) 

1435 

1436 common_ids = current_venue_ids_set.intersection(known_venue_ids_set) 

1437 

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"] 

1450 

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 

1466 

1467 

1468def is_a_valid_row(row: Dict[str, str]) -> bool: 

1469 """ 

1470 This method discards invalid rows in the input CSV file. 

1471 

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 

1536 

1537 

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