Skip to content

Commit 3241a82

Browse files
Merge pull request #154 from wfcommons/snakemake_logger
Snakemake log parser fixes
2 parents 5e56c68 + f96dbb3 commit 3241a82

1 file changed

Lines changed: 36 additions & 7 deletions

File tree

wfcommons/wfinstances/logs/snakemake.py

Lines changed: 36 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -72,6 +72,7 @@ def __init__(self,
7272
self.file_map = {}
7373
self.file_objects = {}
7474
self.task_map = {}
75+
self.task_shell = {}
7576
self.task_input_files = {}
7677
self.task_output_files = {}
7778
self.file_input_output = {}
@@ -98,7 +99,7 @@ def build_workflow(self, workflow_name: Optional[str] = None) -> Workflow:
9899
runtime_system_name=self.wms_name,
99100
runtime_system_url=self.wms_url)
100101

101-
# Parse the sqlite db for to identify task
102+
# Parse the sqlite db for to identify rules
102103
self._build_task_map()
103104

104105
# Parse the sqlite db for to identify files
@@ -122,16 +123,35 @@ def build_workflow(self, workflow_name: Optional[str] = None) -> Workflow:
122123
def _build_task_map(self):
123124
conn = sqlite3.connect(self.snkmt_db)
124125
cursor = conn.cursor()
126+
# Deal with rules
127+
rules = {}
125128
cursor.execute("SELECT * FROM rules")
126129
rows = cursor.fetchall()
127130
for row in rows:
128-
idx = row[0]
129-
task_name = row[1]
130-
if task_name in self.rules_to_ignore:
131+
rule_idx = row[0]
132+
rule_name = row[1]
133+
if rule_name in self.rules_to_ignore:
131134
continue
132-
self.task_map[idx] = task_name
133-
self.task_input_files[idx] = []
134-
self.task_output_files[idx] = []
135+
rules[rule_idx] = rule_name
136+
137+
# Deal with tasks
138+
cursor.execute("SELECT * FROM jobs")
139+
rows = cursor.fetchall()
140+
for row in rows:
141+
task_idx = row[0]
142+
rule_idx = row[3]
143+
# Shell command
144+
if row[8]:
145+
command_list = [x.rstrip().lstrip() for x in row[8].lstrip().rstrip().split('\n')]
146+
shell_cmd = "; ".join(command_list)
147+
else:
148+
shell_cmd = None
149+
if rule_idx not in rules:
150+
continue
151+
self.task_map[task_idx] = rules[rule_idx] + "_" + str(task_idx)
152+
self.task_shell[task_idx] = shell_cmd
153+
self.task_input_files[task_idx] = []
154+
self.task_output_files[task_idx] = []
135155

136156
def _build_file_map(self):
137157
conn = sqlite3.connect(self.snkmt_db)
@@ -181,13 +201,22 @@ def _create_tasks(self):
181201
input_files = [self.file_objects[path] for path in self.task_input_files[idx]]
182202
output_files = [self.file_objects[path] for path in self.task_output_files[idx]]
183203

204+
if self.task_shell[idx]:
205+
program_name = self.task_shell[idx].split(' ')[0]
206+
program_args = self.task_shell[idx].split(' ')[0:]
207+
else:
208+
program_name = "n/a"
209+
program_args = []
210+
184211
task = Task(name=self.task_map[idx],
185212
task_id=self.task_map[idx],
186213
task_type=TaskType.COMPUTE,
187214
runtime=elapsed,
188215
executed_at=start_date,
189216
input_files=input_files,
190217
output_files=output_files,
218+
program=program_name,
219+
args=program_args,
191220
logger=self.logger)
192221
self.workflow.add_task(task)
193222

0 commit comments

Comments
 (0)