|
2 | 2 |
|
3 | 3 | import argparse |
4 | 4 | import re |
| 5 | +import shlex |
5 | 6 | import shutil |
6 | 7 | import subprocess |
7 | 8 | import sys |
|
10 | 11 |
|
11 | 12 |
|
12 | 13 | RMLSTREAMER_JAR_CONTAINER = "/opt/rmlstreamer/RMLStreamer-v2.5.0-standalone.jar" |
| 14 | +_COMMAND_LOGGER = None |
| 15 | + |
| 16 | + |
| 17 | +class CommandLogger: |
| 18 | + def __init__(self, path: Path): |
| 19 | + self.path = path |
| 20 | + self.path.parent.mkdir(parents=True, exist_ok=True) |
| 21 | + self._handle = self.path.open("a", encoding="utf-8") |
| 22 | + |
| 23 | + def run(self, cmd, cwd=None, env=None): |
| 24 | + timestamp = datetime.now().strftime("%Y-%m-%dT%H:%M:%S") |
| 25 | + rendered = " ".join(shlex.quote(str(part)) for part in cmd) |
| 26 | + self._handle.write(f"\n[{timestamp}] $ {rendered}\n") |
| 27 | + if cwd is not None: |
| 28 | + self._handle.write(f"cwd={cwd}\n") |
| 29 | + self._handle.flush() |
| 30 | + |
| 31 | + result = subprocess.run( |
| 32 | + cmd, |
| 33 | + cwd=cwd, |
| 34 | + env=env, |
| 35 | + stdout=self._handle, |
| 36 | + stderr=self._handle, |
| 37 | + text=True, |
| 38 | + ) |
| 39 | + self._handle.write(f"[exit {result.returncode}]\n") |
| 40 | + self._handle.flush() |
| 41 | + return result.returncode |
| 42 | + |
| 43 | + def close(self): |
| 44 | + if not self._handle.closed: |
| 45 | + self._handle.close() |
13 | 46 |
|
14 | 47 |
|
15 | 48 | def eprint(*args): |
16 | 49 | print(*args, file=sys.stderr) |
17 | 50 |
|
18 | 51 |
|
19 | 52 | def run(cmd, cwd=None, env=None): |
20 | | - return subprocess.run(cmd, cwd=cwd, env=env).returncode |
| 53 | + if _COMMAND_LOGGER is not None: |
| 54 | + return _COMMAND_LOGGER.run(cmd, cwd=cwd, env=env) |
| 55 | + return subprocess.run(cmd, cwd=cwd, env=env, capture_output=True, text=True).returncode |
21 | 56 |
|
22 | 57 |
|
23 | 58 | def check_docker(): |
@@ -215,166 +250,179 @@ def main(): |
215 | 250 |
|
216 | 251 | run_id = datetime.now().strftime("%Y%m%dT%H%M%S") |
217 | 252 | timestamp = datetime.now().strftime("%Y-%m-%dT%H:%M:%S") |
218 | | - |
219 | | - if not check_docker(): |
220 | | - return 2 |
| 253 | + wrapper_log_path = metrics_dir / ".wrapper_logs" / f"wrapper-{run_id}.log" |
| 254 | + global _COMMAND_LOGGER |
| 255 | + _COMMAND_LOGGER = CommandLogger(wrapper_log_path) |
| 256 | + print(f"Detailed logs: {wrapper_log_path}") |
221 | 257 |
|
222 | 258 | try: |
223 | | - image_ref, version_requested = resolve_image_ref(args.image, args.image_version) |
224 | | - except ValueError as exc: |
225 | | - eprint(f"Error: {exc}") |
226 | | - return 2 |
| 259 | + if not check_docker(): |
| 260 | + eprint(f"See log for details: {wrapper_log_path}") |
| 261 | + return 2 |
| 262 | + |
| 263 | + try: |
| 264 | + image_ref, version_requested = resolve_image_ref(args.image, args.image_version) |
| 265 | + except ValueError as exc: |
| 266 | + eprint(f"Error: {exc}") |
| 267 | + return 2 |
227 | 268 |
|
228 | | - print("Step 2/5: Ensuring Docker image is available") |
229 | | - if args.build: |
230 | | - if docker_build_image(image_ref, repo_root) != 0: |
231 | | - eprint("Error: docker build failed.") |
| 269 | + print("Step 2/5: Ensuring Docker image is available") |
| 270 | + if args.build: |
| 271 | + print(" - Building Docker image") |
| 272 | + if docker_build_image(image_ref, repo_root) != 0: |
| 273 | + eprint(f"Error: docker build failed. See log: {wrapper_log_path}") |
| 274 | + return 1 |
| 275 | + else: |
| 276 | + if not docker_image_exists(image_ref): |
| 277 | + if version_requested: |
| 278 | + print(f" - Pulling image: {image_ref}") |
| 279 | + if docker_pull_image(image_ref) != 0: |
| 280 | + eprint(f"Error: image version '{image_ref}' not found. See log: {wrapper_log_path}") |
| 281 | + return 2 |
| 282 | + else: |
| 283 | + if args.no_build: |
| 284 | + eprint(f"Error: image '{image_ref}' not found and --no-build set.") |
| 285 | + return 2 |
| 286 | + print(" - Image missing locally, building") |
| 287 | + if docker_build_image(image_ref, repo_root) != 0: |
| 288 | + eprint(f"Error: docker build failed. See log: {wrapper_log_path}") |
| 289 | + return 1 |
| 290 | + |
| 291 | + print("Step 3/5: Converting VCF to TSV") |
| 292 | + tsv_existed = tsv_dir.exists() |
| 293 | + ensure_dir(tsv_dir) |
| 294 | + tsv_cmd = [ |
| 295 | + "sudo", |
| 296 | + "docker", |
| 297 | + "run", |
| 298 | + "--rm", |
| 299 | + "-v", |
| 300 | + f"{str(input_dir)}:/data/in:ro", |
| 301 | + "-v", |
| 302 | + f"{str(tsv_dir)}:/data/tsv", |
| 303 | + image_ref, |
| 304 | + "/opt/vcf-rdfizer/vcf_as_tsv.sh", |
| 305 | + container_input, |
| 306 | + "/data/tsv", |
| 307 | + ] |
| 308 | + if run(tsv_cmd) != 0: |
| 309 | + eprint(f"Error: TSV conversion failed. See log: {wrapper_log_path}") |
232 | 310 | return 1 |
233 | | - else: |
234 | | - if not docker_image_exists(image_ref): |
235 | | - if version_requested: |
236 | | - print(f"Image {image_ref} not found locally. Attempting to pull...") |
237 | | - if docker_pull_image(image_ref) != 0: |
238 | | - eprint(f"Error: image version '{image_ref}' not found.") |
239 | | - return 2 |
240 | | - else: |
241 | | - if args.no_build: |
242 | | - eprint(f"Error: image '{image_ref}' not found and --no-build set.") |
243 | | - return 2 |
244 | | - if docker_build_image(image_ref, repo_root) != 0: |
245 | | - eprint("Error: docker build failed.") |
246 | | - return 1 |
247 | | - |
248 | | - print("Step 3/5: Converting VCF to TSV") |
249 | | - tsv_existed = tsv_dir.exists() |
250 | | - ensure_dir(tsv_dir) |
251 | | - tsv_cmd = [ |
252 | | - "sudo", |
253 | | - "docker", |
254 | | - "run", |
255 | | - "--rm", |
256 | | - "-v", |
257 | | - f"{str(input_dir)}:/data/in:ro", |
258 | | - "-v", |
259 | | - f"{str(tsv_dir)}:/data/tsv", |
260 | | - image_ref, |
261 | | - "/opt/vcf-rdfizer/vcf_as_tsv.sh", |
262 | | - container_input, |
263 | | - "/data/tsv", |
264 | | - ] |
265 | | - if run(tsv_cmd) != 0: |
266 | | - eprint("Error: TSV conversion failed.") |
267 | | - return 1 |
268 | | - |
269 | | - print("Step 4/5: Running Conversion with RMLStreamer") |
270 | | - ensure_dir(out_dir) |
271 | | - ensure_dir(metrics_dir) |
272 | 311 |
|
273 | | - try: |
274 | | - tsv_triplets = discover_tsv_triplets(tsv_dir) |
275 | | - except ValueError as exc: |
276 | | - eprint(f"Error: {exc}") |
277 | | - return 1 |
278 | | - |
279 | | - generated_rules_dir = metrics_dir / "_generated_rules" |
280 | | - if generated_rules_dir.exists(): |
281 | | - shutil.rmtree(generated_rules_dir, ignore_errors=True) |
282 | | - ensure_dir(generated_rules_dir) |
283 | | - |
284 | | - conversion_output_names = [] |
285 | | - for triplet in tsv_triplets: |
286 | | - prefix = triplet["prefix"] |
287 | | - safe_prefix = slugify(prefix) |
288 | | - generated_rules = generated_rules_dir / f"{safe_prefix}.rules.ttl" |
289 | | - render_rules_for_triplet( |
290 | | - rules_path, |
291 | | - generated_rules, |
292 | | - triplet["records"].name, |
293 | | - triplet["headers"].name, |
294 | | - triplet["metadata"].name, |
295 | | - ) |
| 312 | + print("Step 4/5: Running Conversion with RMLStreamer") |
| 313 | + ensure_dir(out_dir) |
| 314 | + ensure_dir(metrics_dir) |
296 | 315 |
|
297 | | - output_name = safe_prefix or slugify(args.out_name) |
298 | | - conversion_output_names.append(output_name) |
299 | | - container_generated_rules = f"/data/rules/{generated_rules.name}" |
| 316 | + try: |
| 317 | + tsv_triplets = discover_tsv_triplets(tsv_dir) |
| 318 | + except ValueError as exc: |
| 319 | + eprint(f"Error: {exc}") |
| 320 | + eprint(f"See log for details: {wrapper_log_path}") |
| 321 | + return 1 |
300 | 322 |
|
301 | | - print(f" - Converting '{prefix}' -> out/{output_name}") |
302 | | - run_cmd = [ |
| 323 | + generated_rules_dir = metrics_dir / "_generated_rules" |
| 324 | + if generated_rules_dir.exists(): |
| 325 | + shutil.rmtree(generated_rules_dir, ignore_errors=True) |
| 326 | + ensure_dir(generated_rules_dir) |
| 327 | + |
| 328 | + conversion_output_names = [] |
| 329 | + for triplet in tsv_triplets: |
| 330 | + prefix = triplet["prefix"] |
| 331 | + safe_prefix = slugify(prefix) |
| 332 | + generated_rules = generated_rules_dir / f"{safe_prefix}.rules.ttl" |
| 333 | + render_rules_for_triplet( |
| 334 | + rules_path, |
| 335 | + generated_rules, |
| 336 | + triplet["records"].name, |
| 337 | + triplet["headers"].name, |
| 338 | + triplet["metadata"].name, |
| 339 | + ) |
| 340 | + |
| 341 | + output_name = safe_prefix or slugify(args.out_name) |
| 342 | + conversion_output_names.append(output_name) |
| 343 | + container_generated_rules = f"/data/rules/{generated_rules.name}" |
| 344 | + |
| 345 | + print(f" - Converting '{prefix}'") |
| 346 | + run_cmd = [ |
| 347 | + "sudo", |
| 348 | + "docker", |
| 349 | + "run", |
| 350 | + "--rm", |
| 351 | + "-v", |
| 352 | + f"{str(generated_rules_dir)}:/data/rules:ro", |
| 353 | + "-v", |
| 354 | + f"{str(tsv_dir)}:/data/tsv:ro", |
| 355 | + "-v", |
| 356 | + f"{str(out_dir)}:/data/out", |
| 357 | + "-v", |
| 358 | + f"{str(metrics_dir)}:/data/metrics", |
| 359 | + "-w", |
| 360 | + "/data/rules", |
| 361 | + "-e", |
| 362 | + f"JAR={RMLSTREAMER_JAR_CONTAINER}", |
| 363 | + "-e", |
| 364 | + f"IN={container_generated_rules}", |
| 365 | + "-e", |
| 366 | + "OUT_DIR=/data/out", |
| 367 | + "-e", |
| 368 | + f"OUT_NAME={output_name}", |
| 369 | + "-e", |
| 370 | + f"RUN_ID={run_id}", |
| 371 | + "-e", |
| 372 | + f"TIMESTAMP={timestamp}", |
| 373 | + "-e", |
| 374 | + f"IN_VCF={container_input}", |
| 375 | + "-e", |
| 376 | + "LOGDIR=/data/metrics", |
| 377 | + image_ref, |
| 378 | + "/opt/vcf-rdfizer/run_conversion.sh", |
| 379 | + ] |
| 380 | + if run(run_cmd) != 0: |
| 381 | + eprint(f"Error: RMLStreamer step failed for '{prefix}'. See log: {wrapper_log_path}") |
| 382 | + return 1 |
| 383 | + |
| 384 | + print("Step 5/5: Compressing outputs") |
| 385 | + compression_out_name = conversion_output_names[0] if len(conversion_output_names) == 1 else "" |
| 386 | + compression_cmd = [ |
303 | 387 | "sudo", |
304 | 388 | "docker", |
305 | 389 | "run", |
306 | 390 | "--rm", |
307 | 391 | "-v", |
308 | | - f"{str(generated_rules_dir)}:/data/rules:ro", |
309 | | - "-v", |
310 | | - f"{str(tsv_dir)}:/data/tsv:ro", |
311 | | - "-v", |
312 | 392 | f"{str(out_dir)}:/data/out", |
313 | 393 | "-v", |
314 | 394 | f"{str(metrics_dir)}:/data/metrics", |
315 | | - "-w", |
316 | | - "/data/rules", |
317 | | - "-e", |
318 | | - f"JAR={RMLSTREAMER_JAR_CONTAINER}", |
319 | 395 | "-e", |
320 | | - f"IN={container_generated_rules}", |
| 396 | + "OUT_ROOT_DIR=/data/out", |
321 | 397 | "-e", |
322 | | - "OUT_DIR=/data/out", |
| 398 | + f"OUT_NAME={compression_out_name}", |
323 | 399 | "-e", |
324 | | - f"OUT_NAME={output_name}", |
| 400 | + "LOGDIR=/data/metrics", |
325 | 401 | "-e", |
326 | 402 | f"RUN_ID={run_id}", |
327 | 403 | "-e", |
328 | 404 | f"TIMESTAMP={timestamp}", |
329 | | - "-e", |
330 | | - f"IN_VCF={container_input}", |
331 | | - "-e", |
332 | | - "LOGDIR=/data/metrics", |
333 | 405 | image_ref, |
334 | | - "/opt/vcf-rdfizer/run_conversion.sh", |
| 406 | + "/opt/vcf-rdfizer/compression.sh", |
| 407 | + "-m", |
| 408 | + args.compression, |
335 | 409 | ] |
336 | | - if run(run_cmd) != 0: |
337 | | - eprint(f"Error: RMLStreamer step failed for '{prefix}'.") |
| 410 | + if run(compression_cmd) != 0: |
| 411 | + eprint(f"Error: compression step failed. See log: {wrapper_log_path}") |
338 | 412 | return 1 |
339 | 413 |
|
340 | | - print("Step 5/5: Compressing outputs") |
341 | | - compression_out_name = conversion_output_names[0] if len(conversion_output_names) == 1 else "" |
342 | | - compression_cmd = [ |
343 | | - "sudo", |
344 | | - "docker", |
345 | | - "run", |
346 | | - "--rm", |
347 | | - "-v", |
348 | | - f"{str(out_dir)}:/data/out", |
349 | | - "-v", |
350 | | - f"{str(metrics_dir)}:/data/metrics", |
351 | | - "-e", |
352 | | - "OUT_ROOT_DIR=/data/out", |
353 | | - "-e", |
354 | | - f"OUT_NAME={compression_out_name}", |
355 | | - "-e", |
356 | | - "LOGDIR=/data/metrics", |
357 | | - "-e", |
358 | | - f"RUN_ID={run_id}", |
359 | | - "-e", |
360 | | - f"TIMESTAMP={timestamp}", |
361 | | - image_ref, |
362 | | - "/opt/vcf-rdfizer/compression.sh", |
363 | | - "-m", |
364 | | - args.compression, |
365 | | - ] |
366 | | - if run(compression_cmd) != 0: |
367 | | - eprint("Error: compression step failed.") |
368 | | - return 1 |
369 | | - |
370 | | - if not args.keep_tsv: |
371 | | - if not tsv_existed: |
372 | | - shutil.rmtree(tsv_dir, ignore_errors=True) |
373 | | - else: |
374 | | - print("Note: TSV directory existed; skipping cleanup.") |
375 | | - |
376 | | - print("Done. See output and metrics directories for results and statistics about the conversion process.") |
377 | | - return 0 |
| 414 | + if not args.keep_tsv: |
| 415 | + if not tsv_existed: |
| 416 | + shutil.rmtree(tsv_dir, ignore_errors=True) |
| 417 | + else: |
| 418 | + print("Note: TSV directory existed; skipping cleanup.") |
| 419 | + |
| 420 | + print("Done. See output and metrics directories for results.") |
| 421 | + return 0 |
| 422 | + finally: |
| 423 | + if _COMMAND_LOGGER is not None: |
| 424 | + _COMMAND_LOGGER.close() |
| 425 | + _COMMAND_LOGGER = None |
378 | 426 |
|
379 | 427 |
|
380 | 428 | if __name__ == "__main__": |
|
0 commit comments