Source code for wfcommons.wfbench.translator.pegasus
#!/usr/bin/env python
# -*- coding: utf-8 -*-
#
# Copyright (c) 2021-2025 The WfCommons Team.
#
# This program is free software: you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation, either version 3 of the License, or
# (at your option) any later version.
import pathlib
from logging import Logger
from typing import Dict, Optional, Union
from .abstract_translator import Translator
from ...common import Workflow
this_dir = pathlib.Path(__file__).resolve().parent
[docs]
class PegasusTranslator(Translator):
"""
A WfFormat parser for creating Pegasus workflow applications.
:param workflow: Workflow benchmark object or path to the workflow benchmark JSON instance.
:type workflow: Union[Workflow, pathlib.Path],
:param logger: The logger where to log information/warning or errors (optional).
:type logger: Logger
"""
def __init__(self,
workflow: Union[Workflow, pathlib.Path],
logger: Optional[Logger] = None) -> None:
"""Create an object of the translator."""
super().__init__(workflow, logger)
self.parsed_tasks = []
self.tasks_map = {}
self.task_counter = 1
[docs]
def translate(self, output_folder: pathlib.Path, tasks_priorities: Optional[Dict[str, int]] = None) -> None:
"""
Translate a workflow benchmark description (WfFormat) into a Pegasus workflow application.
:param output_folder: The path to the folder in which the workflow benchmark will be generated.
:type output_folder: pathlib.Path
:param tasks_priorities: Priorities to be assigned to tasks.
:type tasks_priorities: Optional[Dict[str, int]]
"""
# overall workflow
self.script = f"wf = Workflow('{self.workflow.name}', infer_dependencies=True)\n\n"
# transformation catalog
transformations = []
# tasks' programs
for task in self.tasks.values():
if task.name not in transformations:
transformations.append(task.name)
self.script += "if transformation_path is None:\n" \
f" raise RuntimeError('Unable to find {task.program}')\n" \
f"transformation = Transformation('{task.name}', site='local',\n" \
f" pfn=transformation_path,\n" \
" is_stageable=True)\n" \
"transformation.add_env(PATH='/usr/bin:/bin:.')\n" \
"transformation.add_profiles(Namespace.CONDOR, 'request_disk', '10')\n" \
"tc.add_transformations(transformation)\n\n"
# adding tasks
for task_name in self.root_task_names:
self._add_task(task_name, tasks_priorities=tasks_priorities)
# input file
task = self.tasks[task_name]
for file in task.input_files:
self.script += f"in_file_{self.task_counter} = File('{file.file_id}')\n"
self.script += f"rc.add_replica('local', '{file.file_id}', 'file://' + os.getcwd() + " \
f"'/data/{file.file_id}')\n"
self.script += f"{self.tasks_map[task_name]}.add_inputs(in_file_{self.task_counter})\n" \
f"print('Using input data: ' + os.getcwd() + '/data/{file.file_id}')\n"
self.script += "\n"
# write out the workflow
self.script += "wf.add_replica_catalog(rc)\n" \
"wf.add_transformation_catalog(tc)\n" \
f"wf.write('{self.workflow.name}-benchmark-workflow.yml')\n"
run_workflow_code = self._merge_codelines("templates/pegasus_template.py", self.script)
# write benchmark files
output_folder.mkdir(parents=True)
with open(output_folder.joinpath("pegasus_workflow.py"), "w") as fp:
fp.write(run_workflow_code)
# additional files
self._copy_binary_files(output_folder)
self._generate_input_files(output_folder)
[docs]
def _add_task(self, task_name: str, parent_task: Optional[str] = None, tasks_priorities: Optional[Dict[str, int]] = None) -> None:
"""
Add a task and its dependencies to the workflow.
:param task_name: name of the task
:type task_name: str
:param parent_task: name of the parent task
:type parent_task: Optional[str]
:param tasks_priorities: Priorities to be assigned to tasks.
:type tasks_priorities: Optional[Dict[str, int]]
"""
if task_name not in self.parsed_tasks:
task = self.tasks[task_name]
job_name = f"job_{self.task_counter}"
self.script += f"{job_name} = Job('{task.name}', _id='{task_name}')\n" \
f"task_output_files.setdefault('{job_name}', [])\n"
# task priority
if tasks_priorities and task.name in tasks_priorities:
self.script += f"{job_name}.add_condor_profile(priority='{tasks_priorities[task.name]}')\n"
# find children
children = self.task_children[task_name]
# Generate input spec
input_spec = "\"["
for f in task.input_files:
input_spec += f"\\\\\"{f.file_id}\\\\\","
input_spec = input_spec[:-1] + "]\""
# output files
output_spec = "\"{"
for file in task.output_files:
output_spec += f"\\\\\"{file.file_id}\\\\\":{str(file.size)},"
stage_out = "True" if len(children) == 0 else "False"
self.script += f"out_file_{self.task_counter} = File('{file.file_id}')\n" \
f"task_output_files['{job_name}'].append(out_file_{self.task_counter})\n" \
f"{job_name}.add_outputs(out_file_{self.task_counter}, " \
f"stage_out={stage_out}, register_replica={stage_out})\n"
output_spec = output_spec[:-1] + "}\""
# arguments
args = []
for a in task.args:
if "--output-files" in a:
args.append(f"--output-files {output_spec}")
elif "--input-files" in a:
args.append(f"--input-files {input_spec}")
else:
args.append(a)
args = ", ".join(f"'{a}'" for a in args)
self.script += f"{job_name}.add_args({args})\n"
self.script += f"wf.add_jobs({job_name})\n\n"
self.task_counter += 1
self.parsed_tasks.append(task_name)
self.tasks_map[task_name] = job_name
for child_task_name in children:
self._add_task(child_task_name, job_name, tasks_priorities)
if parent_task:
self.script += f"if '{parent_task}' in task_output_files:\n" \
f" for f in task_output_files['{parent_task}']:\n" \
f" {self.tasks_map[task_name]}.add_inputs(f)\n" \
f"wf.add_dependency({self.tasks_map[task_name]}, parents=[{parent_task}])\n\n"