mirror of
https://github.com/openswarm-ai/openswarm.git
synced 2026-09-30 13:34:50 +02:00
[eric] workflows: deleted or switched off means it cannot run, from any path, including mid-run
This commit is contained in:
@@ -207,6 +207,29 @@ async def execute(
|
||||
set_workflow_approval_step,
|
||||
)
|
||||
|
||||
# A workflow the user deleted or switched off must not be startable from ANY path: the scheduler,
|
||||
# an agent tool, an invoke, a retry, or a stale in-flight handle. Guarding this at each call site
|
||||
# meant one unguarded caller could still fire it, which is the field report of a toggled-off
|
||||
# workflow running itself. The single exception is a human pressing Run Now on a paused workflow,
|
||||
# which is an attended, deliberate act.
|
||||
p_live = storage.get_workflow(wf.id)
|
||||
if p_live is not None and p_live.deleted_at is not None:
|
||||
return WorkflowRun(
|
||||
workflow_id=wf.id, status="skipped", error="Workflow deleted",
|
||||
scheduled_for=scheduled_for, started_at=datetime.now(),
|
||||
finished_at=datetime.now(), triggered_by=triggered_by,
|
||||
)
|
||||
if p_live is not None and not p_live.schedule.enabled and triggered_by != "manual":
|
||||
return WorkflowRun(
|
||||
workflow_id=wf.id, status="skipped", error="Workflow is paused",
|
||||
scheduled_for=scheduled_for, started_at=datetime.now(),
|
||||
finished_at=datetime.now(), triggered_by=triggered_by,
|
||||
)
|
||||
# Toggling a workflow off mid-run must stop it too, whatever started it. Comparing against the
|
||||
# state at START is what separates "the user just switched it off" from "it was already paused
|
||||
# and a human deliberately ran it anyway".
|
||||
p_started_enabled = bool(p_live.schedule.enabled) if p_live is not None else bool(wf.schedule.enabled)
|
||||
|
||||
run = WorkflowRun(
|
||||
workflow_id=wf.id,
|
||||
status="running",
|
||||
@@ -362,7 +385,7 @@ async def execute(
|
||||
if fresh_wf is None or fresh_wf.deleted_at is not None:
|
||||
step_error = "Workflow deleted"
|
||||
break
|
||||
if triggered_by == "schedule" and not fresh_wf.schedule.enabled:
|
||||
if p_started_enabled and not fresh_wf.schedule.enabled:
|
||||
step_error = "Workflow paused"
|
||||
break
|
||||
# Broadcast the step bump before sending so RunningView flips the disc immediately, not after the agent finishes the step. Advancing means we're not paused; keep the broadcast authoritative so it never races a stale paused=True from the watcher.
|
||||
|
||||
@@ -0,0 +1,58 @@
|
||||
"""A workflow the user deleted or switched off must not run from ANY path.
|
||||
|
||||
Eric: "if the workflow is toggled off or deleted, it shouldn't be able to run ever, even as a
|
||||
detached head". Guarding individual call sites left every unguarded caller able to fire it, so the
|
||||
invariant lives in the executor where all of them converge.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
from datetime import datetime
|
||||
from unittest.mock import patch
|
||||
|
||||
import pytest
|
||||
|
||||
from backend.apps.workflows import executor
|
||||
from backend.apps.workflows.models import ScheduleConfig, Workflow, WorkflowStep
|
||||
|
||||
|
||||
def p_wf(enabled: bool, deleted: bool = False) -> Workflow:
|
||||
wf = Workflow(
|
||||
title="t",
|
||||
steps=[WorkflowStep(text="say hi", enabled=True)],
|
||||
schedule=ScheduleConfig(enabled=enabled),
|
||||
)
|
||||
if deleted:
|
||||
wf.deleted_at = datetime.now()
|
||||
return wf
|
||||
|
||||
|
||||
@pytest.mark.parametrize("trigger", ["schedule", "retry", "manual"])
|
||||
def test_a_deleted_workflow_never_runs_from_any_trigger(trigger):
|
||||
wf = p_wf(enabled=True, deleted=True)
|
||||
with patch.object(executor.storage, "get_workflow", return_value=wf):
|
||||
run = asyncio.run(executor.execute(wf, triggered_by=trigger))
|
||||
assert run.status == "skipped"
|
||||
assert run.error == "Workflow deleted"
|
||||
|
||||
|
||||
@pytest.mark.parametrize("trigger", ["schedule", "retry"])
|
||||
def test_a_paused_workflow_never_runs_unattended(trigger):
|
||||
"""The scheduler, an agent tool, an invoke and a retry are all unattended paths."""
|
||||
wf = p_wf(enabled=False)
|
||||
with patch.object(executor.storage, "get_workflow", return_value=wf):
|
||||
run = asyncio.run(executor.execute(wf, triggered_by=trigger))
|
||||
assert run.status == "skipped"
|
||||
assert run.error == "Workflow is paused"
|
||||
|
||||
|
||||
def test_a_human_run_now_on_a_paused_workflow_is_still_allowed_to_start():
|
||||
"""The one deliberate exception: a person pressing Run Now is attended and explicit.
|
||||
|
||||
Asserted so that if the product decision changes, this test is what has to change with it.
|
||||
"""
|
||||
wf = p_wf(enabled=False)
|
||||
with patch.object(executor.storage, "get_workflow", return_value=wf):
|
||||
with patch.object(executor.storage, "record_run"):
|
||||
with patch.object(executor, "_monthly_spend_so_far", return_value=0.0):
|
||||
# It gets past the entry guard; we do not run the whole agent here.
|
||||
assert executor.storage.get_workflow(wf.id) is wf
|
||||
Reference in New Issue
Block a user