Coverage for oc_meta / run / migration / rdf_to_nquads.py: 100%

30 statements  

« prev     ^ index     » next       coverage.py v7.13.4, created at 2026-07-25 10:39 +0000

1#!/usr/bin/env python 

2 

3# Copyright 2026 Arcangelo Massari <arcangelo.massari@unibo.it> 

4# SPDX-FileCopyrightText: 2026 Arcangelo Massari <arcangelo.massari@unibo.it> 

5# 

6# SPDX-License-Identifier: ISC 

7 

8import argparse 

9import multiprocessing 

10import os 

11import zipfile 

12from concurrent.futures import ProcessPoolExecutor 

13from functools import partial 

14from pathlib import Path 

15 

16import py7zr 

17from rdflib import Dataset 

18from rich_argparse import RichHelpFormatter 

19 

20from oc_meta.lib.console import console, create_progress 

21from oc_meta.lib.file_manager import collect_zip_files 

22 

23 

24def process_zip_file( 

25 zip_path: Path, output_dir: Path, input_dir_path: Path, compress: bool 

26) -> None: 

27 graph = Dataset(default_union=True) 

28 with zipfile.ZipFile(zip_path, "r") as zf: 

29 json_file = next(name for name in zf.namelist() if name.endswith(".json")) 

30 with zf.open(json_file) as f: 

31 graph.parse(f, format="json-ld") 

32 

33 nquads_output = graph.serialize(format="nquads") 

34 

35 relative_path = zip_path.relative_to(input_dir_path) 

36 output_filename = str(relative_path).replace(os.sep, "-") 

37 output_filename = Path(output_filename).with_suffix(".nq").name 

38 output_nq_path = output_dir / output_filename 

39 

40 with open(output_nq_path, "w", encoding="utf-8") as f: 

41 f.write(nquads_output) 

42 

43 if compress: 

44 output_7z_path = output_nq_path.with_suffix(".nq.7z") 

45 with py7zr.SevenZipFile(output_7z_path, "w") as archive: 

46 archive.write(output_nq_path, output_filename) 

47 output_nq_path.unlink() 

48 

49 

50def main() -> None: # pragma: no cover 

51 parser = argparse.ArgumentParser( 

52 description="Converts JSON-LD files from ZIP archives to N-Quads format.", 

53 formatter_class=RichHelpFormatter, 

54 ) 

55 parser.add_argument( 

56 "input_dir", 

57 type=str, 

58 help="Input directory containing ZIP files (recursive search)", 

59 ) 

60 parser.add_argument( 

61 "output_dir", type=str, help="Output directory for the converted .nq files" 

62 ) 

63 parser.add_argument( 

64 "-m", 

65 "--mode", 

66 type=str, 

67 choices=["all", "data", "prov"], 

68 default="all", 

69 help="Mode: 'all' for all ZIP files (default), 'data' for entity data only, 'prov' for provenance only", 

70 ) 

71 parser.add_argument( 

72 "-w", 

73 "--workers", 

74 type=int, 

75 default=None, 

76 help="Number of worker processes (defaults to CPU count)", 

77 ) 

78 parser.add_argument( 

79 "-c", 

80 "--compress", 

81 action="store_true", 

82 help="Compress output files using 7z format", 

83 ) 

84 args = parser.parse_args() 

85 

86 input_path = Path(args.input_dir).resolve() 

87 output_path = Path(args.output_dir).resolve() 

88 num_workers = args.workers if args.workers else multiprocessing.cpu_count() 

89 

90 output_path.mkdir(parents=True, exist_ok=True) 

91 

92 zip_files = collect_zip_files( 

93 str(input_path), 

94 only_data=args.mode == "data", 

95 only_prov=args.mode == "prov", 

96 ) 

97 total_files = len(zip_files) 

98 

99 mode_labels = {"all": "", "data": "data ", "prov": "provenance "} 

100 console.print( 

101 f"Found {total_files} {mode_labels[args.mode]}ZIP files in {input_path}" 

102 ) 

103 console.print(f"Output directory: {output_path}") 

104 console.print(f"Workers: {num_workers}") 

105 console.print(f"Compression: {'7z' if args.compress else 'none'}") 

106 

107 fail_count = 0 

108 task_func = partial( 

109 process_zip_file, 

110 output_dir=output_path, 

111 input_dir_path=input_path, 

112 compress=args.compress, 

113 ) 

114 

115 # Use forkserver to avoid deadlocks when forking in a multi-threaded environment 

116 ctx = multiprocessing.get_context("forkserver") 

117 with ProcessPoolExecutor(max_workers=num_workers, mp_context=ctx) as executor: 

118 iterator = executor.map(task_func, zip_files) 

119 

120 with create_progress() as progress: 

121 task = progress.add_task("Converting", total=total_files) 

122 while True: 

123 try: 

124 next(iterator) 

125 progress.update(task, advance=1) 

126 except StopIteration: 

127 break 

128 except Exception as e: 

129 console.print(f"[red]Error: {e}[/red]") 

130 fail_count += 1 

131 progress.update(task, advance=1) 

132 

133 console.print() 

134 console.print("Final report") 

135 console.print(f" Success: {total_files - fail_count}") 

136 console.print(f" Failed: {fail_count}") 

137 

138 

139if __name__ == "__main__": # pragma: no cover 

140 main()