mirror of
https://github.com/openswarm-ai/openswarm.git
synced 2026-09-28 04:24:51 +02:00
[eric] workflows: fix mid-run PATCH/DELETE clobber, optimistic-concurrency PATCH, X-out ghost confirm, custom-draft no-orphan, skipped-run toast,
dup-schedule guard, honest on_missed + Q/R warnings + clearer pause-all + Killed-by-restart copy
This commit is contained in:
@@ -37,6 +37,31 @@ def _resolve_allowed_tools(wf: Workflow) -> list[str]:
|
||||
return list(wf.actions.configured_sets)
|
||||
|
||||
|
||||
def _persist_run_fields(wf: Workflow, run_fields: dict, schedule_runs_count_delta: int = 0) -> None:
|
||||
"""Merge run-side fields into the current on-disk workflow.
|
||||
|
||||
The executor holds the `wf` it was launched with; meanwhile the user
|
||||
may have PATCHed unrelated fields (title, schedule, permissions...).
|
||||
Saving our captured `wf` would clobber those edits. Re-read the
|
||||
authoritative record from storage and only mutate the run-side fields
|
||||
we own. If the workflow has been deleted while we ran, silently skip
|
||||
the save so we don't resurrect a deleted record.
|
||||
|
||||
schedule_runs_count_delta is a small int (0 or 1) that we add to the
|
||||
on-disk schedule.runs_count to avoid the same race overwriting an
|
||||
in-flight bump on the user's PATCH path.
|
||||
"""
|
||||
fresh = storage.get_workflow(wf.id)
|
||||
if fresh is None:
|
||||
# Deleted while we ran. Don't resurrect.
|
||||
return
|
||||
for k, v in run_fields.items():
|
||||
setattr(fresh, k, v)
|
||||
if schedule_runs_count_delta:
|
||||
fresh.schedule.runs_count = fresh.schedule.runs_count + schedule_runs_count_delta
|
||||
storage.save_workflow(fresh)
|
||||
|
||||
|
||||
def _monthly_spend_so_far(wf: Workflow) -> float:
|
||||
"""Sum cost_usd across runs of `wf` started in the last 30 days.
|
||||
|
||||
@@ -80,10 +105,11 @@ async def execute(wf: Workflow, triggered_by: str = "schedule", scheduled_for: O
|
||||
run.error = f"Monthly cost cap reached (${spent:.2f} / ${wf.cost_cap_usd_monthly:.2f})"
|
||||
run.finished_at = datetime.now()
|
||||
storage.record_run(run)
|
||||
wf.last_run_at = run.finished_at
|
||||
wf.last_run_status = "skipped"
|
||||
wf.last_run_id = run.id
|
||||
storage.save_workflow(wf)
|
||||
_persist_run_fields(wf, {
|
||||
"last_run_at": run.finished_at,
|
||||
"last_run_status": "skipped",
|
||||
"last_run_id": run.id,
|
||||
})
|
||||
return run
|
||||
|
||||
storage.record_run(run)
|
||||
@@ -100,7 +126,11 @@ async def execute(wf: Workflow, triggered_by: str = "schedule", scheduled_for: O
|
||||
wf.last_run_at = run.started_at
|
||||
wf.last_run_status = "running"
|
||||
wf.last_run_id = run.id
|
||||
storage.save_workflow(wf)
|
||||
_persist_run_fields(wf, {
|
||||
"last_run_at": run.started_at,
|
||||
"last_run_status": "running",
|
||||
"last_run_id": run.id,
|
||||
})
|
||||
|
||||
try:
|
||||
steps = [s.text for s in wf.steps if s.text and s.text.strip()]
|
||||
@@ -158,11 +188,13 @@ async def execute(wf: Workflow, triggered_by: str = "schedule", scheduled_for: O
|
||||
wf.last_run_status = "success"
|
||||
# Bump runs_count for scheduled fires that reached a terminal state
|
||||
# other than "skipped". Manual runs don't count against max_runs.
|
||||
if triggered_by == "schedule" and run.status in ("success", "ran_late", "failure"):
|
||||
wf.schedule.runs_count += 1
|
||||
runs_delta = 1 if (triggered_by == "schedule" and run.status in ("success", "ran_late", "failure")) else 0
|
||||
storage.record_run(run)
|
||||
wf.last_run_at = run.finished_at
|
||||
storage.save_workflow(wf)
|
||||
_persist_run_fields(wf, {
|
||||
"last_run_at": run.finished_at,
|
||||
"last_run_status": wf.last_run_status,
|
||||
}, schedule_runs_count_delta=runs_delta)
|
||||
except Exception as e:
|
||||
logger.exception("Workflow run failed: %s", e)
|
||||
run.status = "failure"
|
||||
@@ -170,7 +202,10 @@ async def execute(wf: Workflow, triggered_by: str = "schedule", scheduled_for: O
|
||||
run.finished_at = datetime.now()
|
||||
storage.record_run(run)
|
||||
wf.last_run_status = "failure"
|
||||
storage.save_workflow(wf)
|
||||
_persist_run_fields(wf, {
|
||||
"last_run_status": "failure",
|
||||
"last_run_at": run.finished_at,
|
||||
})
|
||||
finally:
|
||||
async with _running_lock:
|
||||
_running.pop(wf.id, None)
|
||||
|
||||
@@ -262,7 +262,12 @@ def _mark_stuck_runs_failed() -> None:
|
||||
for wf in storage.list_workflows():
|
||||
for r in storage.list_runs(wf.id, limit=200):
|
||||
if r.status == "running":
|
||||
storage.update_run(r.id, status="failure", error="Killed by restart", finished_at=now)
|
||||
storage.update_run(
|
||||
r.id,
|
||||
status="failure",
|
||||
error="OpenSwarm closed before this run finished.",
|
||||
finished_at=now,
|
||||
)
|
||||
|
||||
|
||||
def reconcile_on_startup() -> None:
|
||||
|
||||
@@ -4,7 +4,7 @@ from contextlib import asynccontextmanager
|
||||
from datetime import datetime
|
||||
from typing import Optional
|
||||
|
||||
from fastapi import HTTPException
|
||||
from fastapi import HTTPException, Header, Request
|
||||
|
||||
from backend.config.Apps import SubApp
|
||||
from backend.apps.workflows.models import (
|
||||
@@ -183,10 +183,32 @@ async def get_workflow_audit(workflow_id: str, limit: int = 50):
|
||||
|
||||
|
||||
@workflows.router.patch("/{workflow_id}")
|
||||
async def update_workflow(workflow_id: str, body: WorkflowUpdate):
|
||||
async def update_workflow(
|
||||
workflow_id: str,
|
||||
body: WorkflowUpdate,
|
||||
if_match: Optional[str] = Header(default=None, alias="If-Match"),
|
||||
):
|
||||
wf = storage.get_workflow(workflow_id)
|
||||
if not wf:
|
||||
raise HTTPException(status_code=404, detail="Workflow not found")
|
||||
# Optimistic concurrency: if the client passed If-Match, verify it
|
||||
# matches the current updated_at. Stale writes (another window or a
|
||||
# mid-edit background fire) get a 409 so the FE can prompt to reload
|
||||
# instead of silently clobbering the other actor's changes. Missing
|
||||
# header = legacy client, allow through (back-compat with the
|
||||
# frontend's pre-409 code path; FE rolls out If-Match immediately).
|
||||
if if_match:
|
||||
current_stamp = wf.updated_at.isoformat() if hasattr(wf.updated_at, "isoformat") else str(wf.updated_at)
|
||||
# Strip quotes a well-behaved HTTP client might add per RFC 7232.
|
||||
if if_match.strip().strip('"') != current_stamp:
|
||||
raise HTTPException(
|
||||
status_code=409,
|
||||
detail={
|
||||
"error": "stale_update",
|
||||
"message": "This workflow changed in another window or by a recent run. Reload and try again.",
|
||||
"current_updated_at": current_stamp,
|
||||
},
|
||||
)
|
||||
before = wf.model_dump(mode="json")
|
||||
data = body.model_dump(exclude_unset=True)
|
||||
for k, v in data.items():
|
||||
@@ -221,15 +243,20 @@ async def run_workflow_now(workflow_id: str):
|
||||
pre_ids = {r.id for r in storage.list_runs(wf.id, limit=10)}
|
||||
asyncio.create_task(executor.execute(wf, triggered_by="manual"))
|
||||
|
||||
# Poll briefly for the newly created run id (anything not already in
|
||||
# the pre-fire snapshot). Falls back to empty if the executor hasn't
|
||||
# written within 250ms — frontend reconciles via WS afterwards.
|
||||
# Poll briefly for the newly created run id. We also surface the
|
||||
# run's status + error string when it lands quickly (e.g. cost-cap
|
||||
# short-circuit, _running collision) so the FE can render a toast
|
||||
# instead of silently switching to History.
|
||||
for _ in range(25):
|
||||
for r in storage.list_runs(wf.id, limit=10):
|
||||
if r.id not in pre_ids and r.triggered_by == "manual":
|
||||
return {"run_id": r.id}
|
||||
return {
|
||||
"run_id": r.id,
|
||||
"status": r.status,
|
||||
"error": r.error,
|
||||
}
|
||||
await asyncio.sleep(0.01)
|
||||
return {"run_id": ""}
|
||||
return {"run_id": "", "status": None, "error": None}
|
||||
|
||||
|
||||
@workflows.router.get("/{workflow_id}/runs")
|
||||
|
||||
Reference in New Issue
Block a user