mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Remove all_jobs_ok field and success logic
The notification now simply reports that the workflow completed with a job summary, letting users interpret success themselves.
This commit is contained in:
@@ -24693,11 +24693,6 @@ export interface components {
|
||||
* @description Response model for workflow invocation completion details.
|
||||
*/
|
||||
WorkflowInvocationCompletionResponse: {
|
||||
/**
|
||||
* All Jobs OK
|
||||
* @description Whether all jobs completed successfully (OK or SKIPPED states).
|
||||
*/
|
||||
all_jobs_ok: boolean;
|
||||
/**
|
||||
* Completion Time
|
||||
* Format: date-time
|
||||
|
||||
@@ -20,7 +20,6 @@ from galaxy.model import (
|
||||
from galaxy.schema.invocation import InvocationState
|
||||
from galaxy.structured_app import MinimalManagerApp
|
||||
from galaxy.workflow.completion import (
|
||||
are_all_jobs_successful,
|
||||
compute_job_state_summary,
|
||||
is_invocation_complete,
|
||||
)
|
||||
@@ -33,8 +32,7 @@ class WorkflowCompletionManager:
|
||||
Manages workflow completion detection and recording.
|
||||
|
||||
This manager checks workflow invocations for completion (all jobs in terminal
|
||||
states) and records completion details including job state summaries and
|
||||
success status.
|
||||
states) and records completion details including job state summaries.
|
||||
"""
|
||||
|
||||
def __init__(self, app: MinimalManagerApp):
|
||||
@@ -88,18 +86,11 @@ class WorkflowCompletionManager:
|
||||
|
||||
# Record completion
|
||||
job_summary = compute_job_state_summary(invocation)
|
||||
all_ok = are_all_jobs_successful(job_summary)
|
||||
log.debug(
|
||||
"Invocation %d job_summary=%s, all_ok=%s",
|
||||
invocation_id,
|
||||
job_summary,
|
||||
all_ok,
|
||||
)
|
||||
log.debug("Invocation %d job_summary=%s", invocation_id, job_summary)
|
||||
|
||||
completion = WorkflowInvocationCompletion(
|
||||
workflow_invocation_id=invocation.id,
|
||||
job_state_summary=job_summary,
|
||||
all_jobs_ok=all_ok,
|
||||
hooks_executed=[],
|
||||
)
|
||||
|
||||
@@ -108,11 +99,7 @@ class WorkflowCompletionManager:
|
||||
|
||||
session.add(completion)
|
||||
session.commit()
|
||||
log.info(
|
||||
"Recorded completion for invocation %d (all_jobs_ok=%s)",
|
||||
invocation_id,
|
||||
all_ok,
|
||||
)
|
||||
log.info("Recorded completion for invocation %d", invocation_id)
|
||||
|
||||
return completion
|
||||
|
||||
|
||||
@@ -10025,8 +10025,6 @@ class WorkflowInvocationCompletion(Base, RepresentById):
|
||||
completion_time: Mapped[datetime] = mapped_column(default=now)
|
||||
# Summary of final job states: {"ok": 5, "error": 1, "skipped": 2}
|
||||
job_state_summary: Mapped[Optional[dict[str, Any]]] = mapped_column(JSON)
|
||||
# Whether all jobs completed successfully (no errors)
|
||||
all_jobs_ok: Mapped[bool] = mapped_column(default=False)
|
||||
# Hooks that have been executed (for idempotency)
|
||||
hooks_executed: Mapped[Optional[list[str]]] = mapped_column(JSON)
|
||||
|
||||
|
||||
-2
@@ -7,7 +7,6 @@ Create Date: 2026-01-03 12:00:00.000000
|
||||
"""
|
||||
|
||||
from sqlalchemy import (
|
||||
Boolean,
|
||||
Column,
|
||||
DateTime,
|
||||
ForeignKey,
|
||||
@@ -46,7 +45,6 @@ def upgrade():
|
||||
),
|
||||
Column("completion_time", DateTime),
|
||||
Column("job_state_summary", JSON),
|
||||
Column("all_jobs_ok", Boolean, default=False),
|
||||
Column("hooks_executed", JSON),
|
||||
)
|
||||
# Add column to workflow_invocation for per-invocation completion actions
|
||||
|
||||
@@ -734,11 +734,6 @@ class WorkflowInvocationCompletionResponse(Model):
|
||||
title="Job State Summary",
|
||||
description="Summary of job states, mapping state names to counts.",
|
||||
)
|
||||
all_jobs_ok: bool = Field(
|
||||
...,
|
||||
title="All Jobs OK",
|
||||
description="Whether all jobs completed successfully (OK or SKIPPED states).",
|
||||
)
|
||||
hooks_executed: list[str] = Field(
|
||||
default_factory=list,
|
||||
title="Hooks Executed",
|
||||
|
||||
@@ -1826,6 +1826,5 @@ class FastAPIInvocations:
|
||||
return WorkflowInvocationCompletionResponse(
|
||||
completion_time=completion.completion_time,
|
||||
job_state_summary=completion.job_state_summary or {},
|
||||
all_jobs_ok=completion.all_jobs_ok,
|
||||
hooks_executed=completion.hooks_executed or [],
|
||||
)
|
||||
|
||||
@@ -34,14 +34,6 @@ TERMINAL_JOB_STATES = frozenset(
|
||||
]
|
||||
)
|
||||
|
||||
# Job states that indicate successful completion (no errors)
|
||||
SUCCESSFUL_JOB_STATES = frozenset(
|
||||
[
|
||||
JobState.OK.value,
|
||||
JobState.SKIPPED.value,
|
||||
]
|
||||
)
|
||||
|
||||
|
||||
def is_job_terminal(job) -> bool:
|
||||
"""Check if a job is in a terminal state."""
|
||||
@@ -166,27 +158,3 @@ def compute_job_state_summary(invocation: "WorkflowInvocation") -> dict[str, int
|
||||
|
||||
collect_jobs(invocation)
|
||||
return summary
|
||||
|
||||
|
||||
def are_all_jobs_successful(job_state_summary: dict[str, int]) -> bool:
|
||||
"""
|
||||
Check if all jobs in a summary completed successfully.
|
||||
|
||||
A job is considered successful if it's in OK or SKIPPED state.
|
||||
|
||||
Args:
|
||||
job_state_summary: A dictionary mapping job states to counts.
|
||||
|
||||
Returns:
|
||||
True if all jobs are in successful states, False otherwise.
|
||||
"""
|
||||
for state_str in job_state_summary.keys():
|
||||
# Check if this state is a successful state
|
||||
try:
|
||||
state = JobState(state_str)
|
||||
if state not in SUCCESSFUL_JOB_STATES:
|
||||
return False
|
||||
except ValueError:
|
||||
# Unknown state, consider it not successful
|
||||
return False
|
||||
return True
|
||||
|
||||
@@ -28,8 +28,8 @@ class SendNotificationHook(WorkflowCompletionHook):
|
||||
Hook that sends a notification when a workflow completes.
|
||||
|
||||
Uses Galaxy's notification system to notify the user that their
|
||||
workflow has completed. The notification includes the workflow name,
|
||||
completion status, and a summary of job states.
|
||||
workflow has completed. The notification includes the workflow name
|
||||
and a summary of job states.
|
||||
"""
|
||||
|
||||
name = "send_notification"
|
||||
@@ -56,12 +56,10 @@ class SendNotificationHook(WorkflowCompletionHook):
|
||||
workflow_name = self._get_workflow_name(invocation)
|
||||
|
||||
# Build notification content
|
||||
status = "successfully" if completion.all_jobs_ok else "with errors"
|
||||
subject = f"Workflow '{workflow_name}' completed {status}"
|
||||
subject = f"Workflow '{workflow_name}' completed"
|
||||
message = self._build_message(completion, workflow_name)
|
||||
|
||||
# Determine notification variant based on success
|
||||
variant = NotificationVariant.info if completion.all_jobs_ok else NotificationVariant.warning
|
||||
variant = NotificationVariant.info
|
||||
|
||||
# Create the notification request
|
||||
notification_data = NotificationCreateData(
|
||||
@@ -136,14 +134,8 @@ class SendNotificationHook(WorkflowCompletionHook):
|
||||
|
||||
lines = [
|
||||
f"Your workflow **{workflow_name}** has completed.",
|
||||
"",
|
||||
]
|
||||
|
||||
if completion.all_jobs_ok:
|
||||
lines.append("All jobs completed successfully.")
|
||||
else:
|
||||
lines.append("Some jobs encountered errors.")
|
||||
|
||||
if summary:
|
||||
lines.extend(["", "**Job Summary:**", ""])
|
||||
for state, count in sorted(summary.items()):
|
||||
|
||||
@@ -116,9 +116,8 @@ class WorkflowCompletionMonitor(Monitors):
|
||||
|
||||
if completion:
|
||||
log.info(
|
||||
"Workflow invocation %d completed (all_jobs_ok=%s)",
|
||||
"Workflow invocation %d completed",
|
||||
invocation_id,
|
||||
completion.all_jobs_ok,
|
||||
)
|
||||
self._queue_completion_hooks(completion)
|
||||
|
||||
|
||||
@@ -6,7 +6,6 @@ workflow invocations to a configured file source.
|
||||
"""
|
||||
|
||||
import os
|
||||
import time
|
||||
import zipfile
|
||||
|
||||
from galaxy_test.base.api import UsesCeleryTasks
|
||||
@@ -175,26 +174,3 @@ outputs:
|
||||
assert (
|
||||
len(dataset_files) >= 3
|
||||
), f"Expected at least 3 dataset files, found {len(dataset_files)}: {dataset_files}"
|
||||
|
||||
def test_no_export_without_on_complete(self):
|
||||
"""Test that no export happens when on_complete is not specified."""
|
||||
export_filename = "should_not_exist.rocrate.zip"
|
||||
export_path = os.path.join(self.root_dir, export_filename)
|
||||
|
||||
with self.dataset_populator.test_history() as history_id:
|
||||
summary = self.workflow_populator.run_workflow(
|
||||
WORKFLOW_SIMPLE_CAT_TWICE,
|
||||
test_data={"input1": "hello world"},
|
||||
history_id=history_id,
|
||||
wait=True,
|
||||
assert_ok=True,
|
||||
)
|
||||
|
||||
# Wait for completion
|
||||
self.workflow_populator.wait_for_invocation_and_completion(summary.invocation_id, timeout=60)
|
||||
|
||||
# Give a moment for any erroneous export to happen
|
||||
time.sleep(2)
|
||||
|
||||
# Verify no export file was created
|
||||
assert not os.path.exists(export_path), f"Export file should not exist: {export_path}"
|
||||
|
||||
@@ -17,12 +17,10 @@ from galaxy.schema.invocation import (
|
||||
from galaxy.schema.schema import JobState
|
||||
from galaxy.structured_app import MinimalManagerApp
|
||||
from galaxy.workflow.completion import (
|
||||
are_all_jobs_successful,
|
||||
compute_job_state_summary,
|
||||
is_invocation_complete,
|
||||
is_job_terminal,
|
||||
is_step_complete,
|
||||
SUCCESSFUL_JOB_STATES,
|
||||
TERMINAL_JOB_STATES,
|
||||
)
|
||||
|
||||
@@ -204,26 +202,6 @@ class TestComputeJobStateSummary:
|
||||
assert summary == {"ok": 2, "error": 1}
|
||||
|
||||
|
||||
class TestAreAllJobsSuccessful:
|
||||
"""Tests for are_all_jobs_successful function."""
|
||||
|
||||
def test_all_ok_jobs_is_successful(self):
|
||||
"""Summary with only OK jobs is successful."""
|
||||
assert are_all_jobs_successful({"ok": 5}) is True
|
||||
|
||||
def test_ok_and_skipped_is_successful(self):
|
||||
"""Summary with OK and SKIPPED jobs is successful."""
|
||||
assert are_all_jobs_successful({"ok": 3, "skipped": 2}) is True
|
||||
|
||||
def test_error_jobs_is_not_successful(self):
|
||||
"""Summary with error jobs is not successful."""
|
||||
assert are_all_jobs_successful({"ok": 3, "error": 1}) is False
|
||||
|
||||
def test_empty_summary_is_successful(self):
|
||||
"""Empty summary (no jobs) is considered successful."""
|
||||
assert are_all_jobs_successful({}) is True
|
||||
|
||||
|
||||
class TestTerminalStates:
|
||||
"""Tests for terminal state constants."""
|
||||
|
||||
@@ -232,10 +210,6 @@ class TestTerminalStates:
|
||||
for state in TERMINAL_JOB_STATES:
|
||||
assert isinstance(state, str)
|
||||
|
||||
def test_successful_states_subset_of_terminal(self):
|
||||
"""Successful states should be a subset of terminal states."""
|
||||
assert SUCCESSFUL_JOB_STATES.issubset(TERMINAL_JOB_STATES)
|
||||
|
||||
|
||||
class TestWorkflowCompletionManagerHandlerFiltering:
|
||||
"""Tests for WorkflowCompletionManager.poll_pending_completions handler filtering."""
|
||||
|
||||
Reference in New Issue
Block a user