Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions docs/ert/reference/workflows/complete_workflows.rst
Original file line number Diff line number Diff line change
Expand Up @@ -111,6 +111,10 @@ The :code:`status` field is one of :code:`success`, :code:`failed` or
jobs that were stopped because the workflow was cancelled are logged at
:code:`INFO` level.

While an experiment is running, the same information is also shown live
in the *Workflows* tab of the run dialog, without having to read the ERT
log.

Workflows hooked in with :code:`HOOK_WORKFLOW` are in addition recorded
alongside the experiment they belong to, in
:code:`<ENSPATH>/experiments/<experiment_id>/workflow_events.jsonl`. That file
Expand Down
50 changes: 43 additions & 7 deletions src/ert/gui/experiments/run_dialog.py
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,7 @@
RunModelErrorEvent,
RunPathCreatedEvent,
StartingTotalRunPathCreationEvent,
WorkflowEvent,
)
from ert.shared.status.utils import (
byte_with_unit,
Expand All @@ -83,6 +84,7 @@
RealizationWidget,
RunpathProgressWidget,
UpdateWidget,
WorkflowLogWidget,
)
from .view.disk_space_widget import MountType

Expand Down Expand Up @@ -386,12 +388,29 @@ def __init__(
def is_experiment_done(self) -> bool:
return self.flag_experiment_done

def _add_dynamic_tab(self, widget: QWidget, label: str) -> int:
# If the currently selected tab is the most recently added dynamic
# tab, keep following along as new tabs are added; otherwise the
# user has navigated elsewhere and we leave their selection alone.
last_index = self._tab_widget.count() - 1
follow_new_tab = (
last_index >= 0
and not isinstance(self._tab_widget.widget(last_index), WorkflowLogWidget)
and self._tab_widget.currentIndex() == last_index
)
tab_index = self._tab_widget.addTab(widget, label)
if follow_new_tab:
self._tab_widget.setCurrentIndex(tab_index)
return tab_index

def _current_tab_changed(self, index: int) -> None:
widget = self._tab_widget.widget(index)
if isinstance(widget, RealizationWidget):
widget.refresh_current_selection()

self.fm_step_frame.setHidden(isinstance(widget, UpdateWidget))
self.fm_step_frame.setHidden(
isinstance(widget, UpdateWidget | WorkflowLogWidget)
)

@Slot(QModelIndex, int, int)
def on_snapshot_new_iteration(
Expand All @@ -413,14 +432,12 @@ def on_snapshot_new_iteration(
widget.itemClicked.connect(self._select_real)
widget.setProperty("identifier", f"tab-iter-{iteration}")
self._select_real(widget._real_list_model.index(0, 0))
tab_index = self._tab_widget.addTab(
self._add_dynamic_tab(
widget,
f"Realizations for iteration {iteration}"
if not self.is_everest
else f"Batch {iteration}...",
)
if self._tab_widget.currentIndex() == self._tab_widget.count() - 2:
self._tab_widget.setCurrentIndex(tab_index)

if self.is_everest:
self._batch_result_types.append(set())
Expand Down Expand Up @@ -461,6 +478,16 @@ def setup_event_monitoring(
if rerun_failed_realizations is False:
self._snapshot_model.reset()
self._tab_widget.clear()
else:
# Other tabs (Realizations/Update) are updated in place for a
# rerun of failed realizations, but the Workflows tab has no
# such per-event update mechanism, so its previous run's rows
# must be cleared explicitly to avoid mixing old and new output.
for i in range(self._tab_widget.count()):
widget = self._tab_widget.widget(i)
if isinstance(widget, WorkflowLogWidget):
widget.clear()
break

self._worker_thread = QThread(parent=self)

Expand Down Expand Up @@ -606,9 +633,7 @@ def _on_event(self, event: object) -> None:
self.progress_update_event.emit(status_count, realization_count)
case RunModelUpdateBeginEvent(iteration=iteration):
widget = UpdateWidget(iteration)
tab_index = self._tab_widget.addTab(widget, f"Update {iteration}")
if self._tab_widget.currentIndex() == self._tab_widget.count() - 2:
self._tab_widget.setCurrentIndex(tab_index)
self._add_dynamic_tab(widget, f"Update {iteration}")
widget.begin(event)
case RunModelUpdateEndEvent():
self._progress_widget.stop_waiting_progress_bar()
Expand All @@ -622,6 +647,8 @@ def _on_event(self, event: object) -> None:
case RunModelErrorEvent():
self._get_update_widget(event.iteration).error(event)
event.write_as_csv(self.output_path)
case WorkflowEvent():
self._get_or_create_workflow_log_widget().add_event(event)
case EverestBatchResultEvent():
batch_types = self._batch_result_types[event.batch]
batch_types.add(event.result_type)
Expand Down Expand Up @@ -660,6 +687,15 @@ def _get_update_widget(self, iteration: int) -> UpdateWidget:
return widget
raise ValueError("Could not find UpdateWidget")

def _get_or_create_workflow_log_widget(self) -> WorkflowLogWidget:
for i in range(self._tab_widget.count()):
widget = self._tab_widget.widget(i)
if isinstance(widget, WorkflowLogWidget):
return widget
workflow_log_widget = WorkflowLogWidget(self)
self._tab_widget.insertTab(0, workflow_log_widget, "Workflows")
return workflow_log_widget

def update_total_progress(
self, progress_value: float, iteration_label: str, iteration: int | None = None
) -> None:
Expand Down
2 changes: 2 additions & 0 deletions src/ert/gui/experiments/view/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,11 +3,13 @@
from .realization import RealizationWidget
from .runpath_progress_widget import RunpathProgressWidget
from .update import UpdateWidget
from .workflow_log import WorkflowLogWidget

__all__ = [
"DiskSpaceWidget",
"ProgressWidget",
"RealizationWidget",
"RunpathProgressWidget",
"UpdateWidget",
"WorkflowLogWidget",
]
226 changes: 226 additions & 0 deletions src/ert/gui/experiments/view/workflow_log.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,226 @@
from __future__ import annotations

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

perhaps name the file workflow_log_widget?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

all the entries in this folder drop the widget suffix, so that is fine.
The question for me is whether to call it rather WorkflowLogInfoWidget -> workflow_log_info since workflow_log might not convey the exact info.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not all entries do. disk_space_widget.py, progress_widget.py ..


from typing import cast

from PyQt6.QtCore import Qt
from PyQt6.QtGui import QColor, QFontDatabase
from PyQt6.QtWidgets import (
QAbstractItemView,
QComboBox,
QHBoxLayout,
QHeaderView,
QLabel,
QPlainTextEdit,
QSplitter,
QTableWidget,
QTableWidgetItem,
QVBoxLayout,
QWidget,
)

from ert.ensemble_evaluator.state import COLOR_CANCELLED, COLOR_FAILED, COLOR_FINISHED
from ert.run_models.event import WorkflowEvent
from ert.workflow_runner import WorkflowJobStatus

NO_ITERATION_LABEL = "Pre/post experiment"
NO_OUTPUT_PLACEHOLDER = "(no output)"
_EXPERIMENT_WIDE_HOOKS = frozenset({"PRE_EXPERIMENT", "POST_EXPERIMENT"})

_COLUMNS = ("Hook", "Workflow", "Job", "Status", "Time")


class WorkflowLogWidget(QWidget):
# Table of workflow job invocations with their captured output.

def __init__(self, parent: QWidget | None = None) -> None:
super().__init__(parent)

self._events: dict[int | None, list[WorkflowEvent]] = {}
self._iteration_chosen_by_user = False

self._iteration_selector = QComboBox(self)
self._iteration_selector.currentIndexChanged.connect(self._on_iteration_changed)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

why do we need a name here?

self._iteration_selector.activated.connect(self._on_iteration_chosen_by_user)

selector_row = QHBoxLayout()
selector_row.setSpacing(6)
selector_row.addWidget(QLabel("Iteration:"))
selector_row.addWidget(self._iteration_selector)
selector_row.addStretch()

self._table = QTableWidget(0, len(_COLUMNS), self)
self._table.setHorizontalHeaderLabels(_COLUMNS)
vertical_header = self._table.verticalHeader()
assert vertical_header is not None
vertical_header.setVisible(False)
self._table.setEditTriggers(QAbstractItemView.EditTrigger.NoEditTriggers)
self._table.setSelectionBehavior(QAbstractItemView.SelectionBehavior.SelectRows)
self._table.setSelectionMode(QAbstractItemView.SelectionMode.SingleSelection)
header = self._table.horizontalHeader()
assert header is not None
header.setSectionResizeMode(QHeaderView.ResizeMode.ResizeToContents)
header.setSectionResizeMode(
_COLUMNS.index("Job"), QHeaderView.ResizeMode.Stretch
)
self._table.itemSelectionChanged.connect(self._on_row_selected)

self._stdout_view = self._make_output_view()
self._stderr_view = self._make_output_view()

detail = QSplitter(Qt.Orientation.Horizontal, self)
detail.addWidget(self._make_labelled_output("Stdout", self._stdout_view))
detail.addWidget(self._make_labelled_output("Stderr", self._stderr_view))

splitter = QSplitter(Qt.Orientation.Vertical, self)
splitter.addWidget(self._table)
splitter.addWidget(detail)
splitter.setStretchFactor(0, 2)
splitter.setStretchFactor(1, 1)

layout = QVBoxLayout(self)
layout.setContentsMargins(8, 8, 8, 8)
layout.setSpacing(8)
layout.addLayout(selector_row)
layout.addWidget(splitter)

self._clear_detail()

def _make_output_view(self) -> QPlainTextEdit:
view = QPlainTextEdit(self)
view.setReadOnly(True)
view.setLineWrapMode(QPlainTextEdit.LineWrapMode.NoWrap)
view.setFont(QFontDatabase.systemFont(QFontDatabase.SystemFont.FixedFont))
return view

def _make_labelled_output(self, title: str, view: QPlainTextEdit) -> QWidget:
container = QWidget(self)
layout = QVBoxLayout(container)
layout.setContentsMargins(0, 0, 0, 0)
layout.setSpacing(2)
layout.addWidget(QLabel(title, container))
layout.addWidget(view)
return container

def add_event(self, event: WorkflowEvent) -> None:
group = self._group_key(event)
is_new_group = group not in self._events
self._events.setdefault(group, []).append(event)

if is_new_group:
self._rebuild_iteration_selector()
elif group == self._selected_iteration():
self._append_row(event)

def _group_key(self, event: WorkflowEvent) -> int | None:
if event.hook in _EXPERIMENT_WIDE_HOOKS:
return None
return event.iteration

def _selected_iteration(self) -> int | None:
index = self._iteration_selector.currentIndex()
if index < 0:
return None
return cast(int | None, self._iteration_selector.itemData(index))

def _rebuild_iteration_selector(self) -> None:
previously_selected = self._selected_iteration()
had_selection = self._iteration_selector.count() > 0

iterations: list[int | None] = []
if None in self._events:
iterations.append(None)
iterations.extend(sorted(i for i in self._events if i is not None))

self._iteration_selector.blockSignals(True)
self._iteration_selector.clear()
for iteration in iterations:
label = (
NO_ITERATION_LABEL if iteration is None else f"Iteration {iteration}"
)
self._iteration_selector.addItem(label, iteration)
if (
self._iteration_chosen_by_user
and had_selection
and previously_selected in iterations
):
self._iteration_selector.setCurrentIndex(
iterations.index(previously_selected)
)
else:
# Follow the newest iteration until the user picks one themselves.
self._iteration_selector.setCurrentIndex(len(iterations) - 1)
self._iteration_selector.blockSignals(False)
self._on_iteration_changed()

def _on_iteration_chosen_by_user(self, _index: int) -> None:
self._iteration_chosen_by_user = True

def _on_iteration_changed(self) -> None:
self._table.clearContents()
self._table.setRowCount(0)
for event in self._events.get(self._selected_iteration(), []):
self._append_row(event)
self._clear_detail()

def _append_row(self, event: WorkflowEvent) -> None:
row = self._table.rowCount()
self._table.insertRow(row)

job = event.job_name
if event.arguments:
job += f"({', '.join(event.arguments)})"
status, color = self._status_and_color(event)

values = (
event.hook,
event.workflow_name,
job,
status,
event.timestamp.strftime("%H:%M:%S"),
)
for column, value in enumerate(values):
item = QTableWidgetItem(value)
if column == 3:
item.setBackground(color)
item.setForeground(QColor(0, 0, 0))
self._table.setItem(row, column, item)

def _status_and_color(self, event: WorkflowEvent) -> tuple[str, QColor]:
match event.status:
case WorkflowJobStatus.CANCELLED:
return "Cancelled", QColor(*COLOR_CANCELLED)
case WorkflowJobStatus.FAILED:
return "Failed", QColor(*COLOR_FAILED)
case WorkflowJobStatus.SUCCESS:
return "Succeeded", QColor(*COLOR_FINISHED)

def _on_row_selected(self) -> None:
selection_model = self._table.selectionModel()
assert selection_model is not None
rows = selection_model.selectedRows()
if not rows:
self._clear_detail()
return
events = self._events.get(self._selected_iteration(), [])
row = rows[0].row()
if row >= len(events):
self._clear_detail()
return
event = events[row]
self._stdout_view.setPlainText(event.stdout or NO_OUTPUT_PLACEHOLDER)
self._stderr_view.setPlainText(event.stderr or NO_OUTPUT_PLACEHOLDER)

def _clear_detail(self) -> None:
self._stdout_view.setPlainText(NO_OUTPUT_PLACEHOLDER)
self._stderr_view.setPlainText(NO_OUTPUT_PLACEHOLDER)

def clear(self) -> None:
"""Discard all displayed workflow events, e.g. before a rerun."""
self._events = {}
self._iteration_chosen_by_user = False
self._iteration_selector.blockSignals(True)
self._iteration_selector.clear()
self._iteration_selector.blockSignals(False)
self._table.clearContents()
self._table.setRowCount(0)
self._clear_detail()
Loading