mirror of
https://github.com/openswarm-ai/openswarm.git
synced 2026-10-01 14:04:51 +02:00
[eric] workflows: hand weekly and self-ending schedules to the cloud, zone and all
This commit is contained in:
@@ -15,7 +15,7 @@ from pydantic import BaseModel, ConfigDict, Field
|
||||
from typeguard import typechecked
|
||||
|
||||
from backend.apps.settings.credentials import account_auth
|
||||
from backend.apps.workflows.cloud.schedule import CloudSchedule
|
||||
from backend.apps.workflows.cloud.schedule import CloudSchedule, wire
|
||||
|
||||
# The cloud router is mounted at /api/workflows and a trailing slash 404s there, so the collection paths are the empty string, not "/".
|
||||
COLLECTION = ""
|
||||
@@ -71,6 +71,9 @@ class HostedWorkflow(BaseModel):
|
||||
id: str
|
||||
enabled: bool = False
|
||||
next_run_at: Optional[int] = None
|
||||
# Total fires, cloud plus the ones done here before the handover. None from a control plane that
|
||||
# predates the count, and a "we were not told" must never be read as a zero.
|
||||
runs_done: Optional[int] = None
|
||||
|
||||
|
||||
class CloudPreflight(BaseModel):
|
||||
@@ -160,10 +163,12 @@ def p_hosted(raw: Any) -> Optional[HostedWorkflow]:
|
||||
if not isinstance(ident, str):
|
||||
return None
|
||||
nxt = raw.get("next_run_at")
|
||||
done = raw.get("runs_done")
|
||||
return HostedWorkflow(
|
||||
id=ident,
|
||||
enabled=bool(raw.get("enabled")),
|
||||
next_run_at=nxt if isinstance(nxt, int) else None,
|
||||
runs_done=done if isinstance(done, int) else None,
|
||||
)
|
||||
|
||||
|
||||
@@ -202,16 +207,13 @@ async def preflight(definition: Dict[str, Any], hosted_id: Optional[str]) -> Clo
|
||||
if not isinstance(raw, dict):
|
||||
raise CloudUnreachable("the cloud sent a preflight we could not read")
|
||||
limits, usage = p_allowance(raw)
|
||||
capability = raw.get("capability")
|
||||
cap = raw.get("capability")
|
||||
return CloudPreflight(
|
||||
plan=raw.get("plan") if isinstance(raw.get("plan"), str) else None,
|
||||
limits=limits,
|
||||
usage=usage,
|
||||
capability=(
|
||||
CloudCapability(ok=bool(capability.get("ok")), reason=capability.get("reason"))
|
||||
if isinstance(capability, dict)
|
||||
else None
|
||||
),
|
||||
capability=CloudCapability(ok=bool(cap.get("ok")), reason=cap.get("reason"))
|
||||
if isinstance(cap, dict) else None,
|
||||
hosted=p_hosted(raw.get("hosted")),
|
||||
)
|
||||
|
||||
@@ -229,11 +231,13 @@ async def p_preflight_from_list(hosted_id: Optional[str]) -> CloudPreflight:
|
||||
|
||||
@typechecked
|
||||
async def put_workflow(
|
||||
*, hosted_id: Optional[str], name: str, definition: Dict[str, Any], schedule: CloudSchedule
|
||||
*, hosted_id: Optional[str], name: str, definition: Dict[str, Any],
|
||||
schedule: CloudSchedule, runs_before: int = 0,
|
||||
) -> HostedWorkflow:
|
||||
"""Create the hosted copy, or re-push onto the existing row so an edited
|
||||
workflow stops running last week's prose."""
|
||||
body = {"name": name, "definition": definition, "schedule": schedule.model_dump()}
|
||||
"""Create the hosted copy, or re-push onto the existing row so an edited workflow stops running
|
||||
last week's prose. runs_before rides only on the create: an edit that resent it would hand a
|
||||
nearly-spent run cap its whole budget back."""
|
||||
body: Dict[str, Any] = {"name": name, "definition": definition, "schedule": wire(schedule)}
|
||||
if hosted_id:
|
||||
try:
|
||||
raw = await p_call("POST", f"/{hosted_id}/update", body)
|
||||
@@ -244,7 +248,7 @@ async def put_workflow(
|
||||
# 404 is the row being gone (deleted elsewhere, or a control plane with no update route); make a fresh one.
|
||||
if exc.status != 404:
|
||||
raise
|
||||
raw = await p_call("POST", COLLECTION, body)
|
||||
raw = await p_call("POST", COLLECTION, {**body, "runs_before": max(0, runs_before)})
|
||||
hosted = p_hosted(raw)
|
||||
if not hosted:
|
||||
raise CloudUnreachable("the cloud accepted the workflow but did not say which one")
|
||||
@@ -281,16 +285,14 @@ async def list_runs(hosted_id: str) -> List[CloudRun]:
|
||||
if not isinstance(row, dict):
|
||||
continue
|
||||
notices = row.get("notices")
|
||||
out.append(
|
||||
CloudRun(
|
||||
id=str(row.get("id") or ""),
|
||||
status=str(row.get("status") or "unknown"),
|
||||
started_at=row.get("started_at") if isinstance(row.get("started_at"), int) else None,
|
||||
finished_at=row.get("finished_at") if isinstance(row.get("finished_at"), int) else None,
|
||||
error=row.get("error") if isinstance(row.get("error"), str) else None,
|
||||
answer=row.get("answer") if isinstance(row.get("answer"), str) else None,
|
||||
notices=[n for n in notices if isinstance(n, str)] if isinstance(notices, list) else [],
|
||||
cost_usd=row.get("cost_usd") if isinstance(row.get("cost_usd"), (int, float)) else None,
|
||||
)
|
||||
)
|
||||
out.append(CloudRun(
|
||||
id=str(row.get("id") or ""),
|
||||
status=str(row.get("status") or "unknown"),
|
||||
started_at=row.get("started_at") if isinstance(row.get("started_at"), int) else None,
|
||||
finished_at=row.get("finished_at") if isinstance(row.get("finished_at"), int) else None,
|
||||
error=row.get("error") if isinstance(row.get("error"), str) else None,
|
||||
answer=row.get("answer") if isinstance(row.get("answer"), str) else None,
|
||||
notices=[n for n in notices if isinstance(n, str)] if isinstance(notices, list) else [],
|
||||
cost_usd=row.get("cost_usd") if isinstance(row.get("cost_usd"), (int, float)) else None,
|
||||
))
|
||||
return out
|
||||
|
||||
@@ -18,7 +18,7 @@ from typeguard import typechecked
|
||||
from backend.apps.workflows import scheduler, storage
|
||||
from backend.apps.workflows.cloud import client as cloud
|
||||
from backend.apps.workflows.cloud.definition import cloud_definition, definition_signature
|
||||
from backend.apps.workflows.cloud.schedule import ScheduleSupported, to_cloud_schedule
|
||||
from backend.apps.workflows.cloud.schedule import ScheduleSupported, to_cloud_schedule, wire
|
||||
from backend.apps.workflows.cloud.status import epoch_to_datetime
|
||||
from backend.apps.workflows.models import Workflow
|
||||
|
||||
@@ -57,6 +57,7 @@ async def hand_to_cloud(wf: Workflow, enabled: bool) -> TargetOutcome:
|
||||
name=wf.title or "Workflow",
|
||||
definition=definition,
|
||||
schedule=mapping.schedule,
|
||||
runs_before=wf.schedule.runs_count,
|
||||
)
|
||||
if hosted.enabled != enabled:
|
||||
hosted = await cloud.set_enabled(hosted.id, enabled)
|
||||
@@ -70,7 +71,7 @@ async def hand_to_cloud(wf: Workflow, enabled: bool) -> TargetOutcome:
|
||||
|
||||
wf.execution_target = "cloud"
|
||||
wf.cloud_workflow_id = hosted.id
|
||||
wf.cloud_definition_signature = definition_signature(definition, mapping.schedule.model_dump())
|
||||
wf.cloud_definition_signature = definition_signature(definition, wire(mapping.schedule))
|
||||
wf.schedule.enabled = enabled
|
||||
wf.next_run_at = epoch_to_datetime(hosted.next_run_at) if enabled else None
|
||||
wf.updated_at = datetime.now()
|
||||
|
||||
@@ -1,16 +1,22 @@
|
||||
"""Map a local schedule onto the cloud scheduler's much smaller vocabulary.
|
||||
"""Map a local schedule onto the cloud scheduler's smaller vocabulary.
|
||||
|
||||
The cloud speaks two cadences: repeat every N minutes, or once a day at a UTC
|
||||
time. Everything else this app can express (set weekdays, monthly, every third
|
||||
day, stop after N runs) has no cloud equivalent, and silently rounding one of
|
||||
them off would fire a workflow on days the user never picked. So anything that
|
||||
does not map exactly is refused here, in the user's own words, and stays on
|
||||
their machine where it already works.
|
||||
The cloud speaks three cadences: repeat every N minutes, once a day, or on the
|
||||
weekdays you picked. Each can carry an end date and a run cap. What it does not
|
||||
speak is a cadence with a phase longer than one period (every third day, every
|
||||
other week, monthly), because the phase is anchored to the workflow's creation
|
||||
on this machine and there is nowhere on the wire to put that anchor. Silently
|
||||
rounding one of those off would fire a workflow on days the user never picked,
|
||||
so it is refused here in the user's own words and stays on their machine.
|
||||
|
||||
The wall-clock kinds carry the IANA zone rather than a UTC hour. A UTC hour is a
|
||||
schedule that moves by an hour twice a year for everyone outside UTC: "9am" set
|
||||
in July quietly becomes 8am in November. The cloud does its recurrence maths in
|
||||
the zone for the same reason scheduler._next_fire_after does.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime
|
||||
from typing import Literal, Optional, Union
|
||||
from typing import Any, Dict, List, Literal, Optional, Union
|
||||
from zoneinfo import ZoneInfo, ZoneInfoNotFoundError
|
||||
|
||||
from pydantic import BaseModel, ConfigDict
|
||||
@@ -19,25 +25,40 @@ from typeguard import typechecked
|
||||
from backend.apps.workflows.models import ScheduleConfig
|
||||
from backend.apps.workflows.scheduler import host_timezone_name
|
||||
|
||||
CADENCE_PREFIX = "Cloud runs repeat on an interval or once a day."
|
||||
CADENCE_PREFIX = "Cloud runs repeat on an interval, once a day, or on the weekdays you pick."
|
||||
|
||||
|
||||
class CloudIntervalSchedule(BaseModel):
|
||||
class CloudScheduleBase(BaseModel):
|
||||
model_config = ConfigDict(validate_assignment=True)
|
||||
|
||||
# Both optional, both meaning "and then it is finished". Sent only when set: the cloud's schema
|
||||
# takes an absent field, not a null one.
|
||||
ends_at: Optional[int] = None
|
||||
max_runs: Optional[int] = None
|
||||
|
||||
|
||||
class CloudIntervalSchedule(CloudScheduleBase):
|
||||
kind: Literal["interval"] = "interval"
|
||||
minutes: int
|
||||
|
||||
|
||||
class CloudDailySchedule(BaseModel):
|
||||
model_config = ConfigDict(validate_assignment=True)
|
||||
|
||||
class CloudDailySchedule(CloudScheduleBase):
|
||||
kind: Literal["daily"] = "daily"
|
||||
hour_utc: int
|
||||
minute_utc: int
|
||||
hour: int
|
||||
minute: int
|
||||
timezone: str
|
||||
|
||||
|
||||
CloudSchedule = Union[CloudIntervalSchedule, CloudDailySchedule]
|
||||
class CloudWeeklySchedule(CloudScheduleBase):
|
||||
kind: Literal["weekly"] = "weekly"
|
||||
# Sunday=0, the same convention as ScheduleConfig.on_days.
|
||||
days: List[int]
|
||||
hour: int
|
||||
minute: int
|
||||
timezone: str
|
||||
|
||||
|
||||
CloudSchedule = Union[CloudIntervalSchedule, CloudDailySchedule, CloudWeeklySchedule]
|
||||
|
||||
|
||||
class ScheduleSupported(BaseModel):
|
||||
@@ -58,62 +79,87 @@ ScheduleMapping = Union[ScheduleSupported, ScheduleUnsupported]
|
||||
|
||||
|
||||
@typechecked
|
||||
def p_zone(name: str) -> ZoneInfo:
|
||||
if not name or name == "local":
|
||||
name = host_timezone_name()
|
||||
try:
|
||||
return ZoneInfo(name)
|
||||
except ZoneInfoNotFoundError:
|
||||
return ZoneInfo("UTC")
|
||||
def wire(sched: CloudSchedule) -> Dict[str, Any]:
|
||||
"""The JSON body shape. Unset bounds are dropped rather than sent as null, which is what the
|
||||
cloud's schema expects and what keeps the definition fingerprint stable across versions."""
|
||||
return sched.model_dump(exclude_none=True)
|
||||
|
||||
|
||||
@typechecked
|
||||
def p_utc_time_of_day(sched: ScheduleConfig, ref: Optional[datetime] = None) -> CloudDailySchedule:
|
||||
"""The user's wall-clock time expressed in UTC, using today's offset."""
|
||||
zone = p_zone(sched.timezone)
|
||||
local = (ref or datetime.now(zone)).astimezone(zone)
|
||||
at = local.replace(hour=sched.hour, minute=sched.minute, second=0, microsecond=0)
|
||||
utc = at.astimezone(ZoneInfo("UTC"))
|
||||
return CloudDailySchedule(hour_utc=utc.hour, minute_utc=utc.minute)
|
||||
def p_zone_name(name: str) -> str:
|
||||
"""A concrete IANA name the cloud can hand to its own tz database. "local" and anything
|
||||
unresolvable fall back to this host's zone, which is what the local scheduler already does."""
|
||||
if not name or name == "local":
|
||||
return host_timezone_name()
|
||||
try:
|
||||
ZoneInfo(name)
|
||||
except ZoneInfoNotFoundError:
|
||||
return host_timezone_name()
|
||||
return name
|
||||
|
||||
|
||||
@typechecked
|
||||
def p_epoch_ms(when: datetime) -> int:
|
||||
# Naive datetimes are host-local, matching how the local scheduler reads its own stored dates.
|
||||
aware = when if when.tzinfo is not None else when.replace(tzinfo=ZoneInfo(host_timezone_name()))
|
||||
return int(aware.timestamp() * 1000)
|
||||
|
||||
|
||||
@typechecked
|
||||
def p_bounds(sched: ScheduleConfig) -> Dict[str, Any]:
|
||||
out: Dict[str, Any] = {}
|
||||
if sched.ends_at is not None:
|
||||
out["ends_at"] = p_epoch_ms(sched.ends_at)
|
||||
if sched.max_runs is not None:
|
||||
out["max_runs"] = sched.max_runs
|
||||
return out
|
||||
|
||||
|
||||
@typechecked
|
||||
def to_cloud_schedule(sched: ScheduleConfig) -> ScheduleMapping:
|
||||
if sched.max_runs is not None:
|
||||
return ScheduleUnsupported(
|
||||
reason=(
|
||||
f"This schedule stops itself after {sched.max_runs} "
|
||||
f"run{'' if sched.max_runs == 1 else 's'}, and the cloud scheduler cannot count down "
|
||||
"to a stop. Remove the limit to run it in the cloud."
|
||||
),
|
||||
)
|
||||
if sched.ends_at is not None:
|
||||
return ScheduleUnsupported(
|
||||
reason=(
|
||||
"This schedule has an end date, and the cloud scheduler cannot honour one. "
|
||||
"Remove the end date to run it in the cloud."
|
||||
),
|
||||
)
|
||||
bounds = p_bounds(sched)
|
||||
if sched.repeat_unit == "minute":
|
||||
return ScheduleSupported(schedule=CloudIntervalSchedule(minutes=max(5, sched.repeat_every)))
|
||||
return ScheduleSupported(schedule=CloudIntervalSchedule(minutes=max(5, sched.repeat_every), **bounds))
|
||||
if sched.repeat_unit == "hour":
|
||||
return ScheduleSupported(schedule=CloudIntervalSchedule(minutes=max(5, sched.repeat_every * 60)))
|
||||
if sched.repeat_unit == "day" and sched.repeat_every == 1:
|
||||
return ScheduleSupported(schedule=p_utc_time_of_day(sched))
|
||||
return ScheduleSupported(
|
||||
schedule=CloudIntervalSchedule(minutes=max(5, sched.repeat_every * 60), **bounds)
|
||||
)
|
||||
|
||||
zone = p_zone_name(sched.timezone)
|
||||
if sched.repeat_unit == "day":
|
||||
if sched.repeat_every == 1:
|
||||
return ScheduleSupported(
|
||||
schedule=CloudDailySchedule(hour=sched.hour, minute=sched.minute, timezone=zone, **bounds)
|
||||
)
|
||||
return ScheduleUnsupported(
|
||||
reason=(
|
||||
f"{CADENCE_PREFIX} This one runs every {sched.repeat_every} days at a set time, "
|
||||
"which the cloud scheduler cannot do yet, so it stays on this device."
|
||||
),
|
||||
)
|
||||
|
||||
if sched.repeat_unit == "week":
|
||||
if not sched.on_days:
|
||||
return ScheduleUnsupported(
|
||||
reason="Pick the days this should run on before choosing where it runs.",
|
||||
)
|
||||
if sched.repeat_every == 1:
|
||||
return ScheduleSupported(
|
||||
schedule=CloudWeeklySchedule(
|
||||
days=sorted(sched.on_days),
|
||||
hour=sched.hour,
|
||||
minute=sched.minute,
|
||||
timezone=zone,
|
||||
**bounds,
|
||||
)
|
||||
)
|
||||
return ScheduleUnsupported(
|
||||
reason=(
|
||||
f"{CADENCE_PREFIX} This one runs on the weekdays you picked, "
|
||||
f"{CADENCE_PREFIX} This one runs every {sched.repeat_every} weeks, "
|
||||
"which the cloud scheduler cannot do yet, so it stays on this device."
|
||||
),
|
||||
)
|
||||
|
||||
return ScheduleUnsupported(
|
||||
reason=(
|
||||
f"{CADENCE_PREFIX} This one runs monthly, "
|
||||
|
||||
@@ -16,7 +16,7 @@ from typeguard import typechecked
|
||||
from backend.apps.workflows import storage
|
||||
from backend.apps.workflows.cloud import client as cloud
|
||||
from backend.apps.workflows.cloud.definition import cloud_definition, definition_signature
|
||||
from backend.apps.workflows.cloud.schedule import ScheduleSupported, to_cloud_schedule
|
||||
from backend.apps.workflows.cloud.schedule import ScheduleSupported, to_cloud_schedule, wire
|
||||
from backend.apps.workflows.models import Workflow
|
||||
|
||||
|
||||
@@ -74,7 +74,7 @@ def current_signature(wf: Workflow) -> Optional[str]:
|
||||
mapping = to_cloud_schedule(wf.schedule)
|
||||
if not isinstance(mapping, ScheduleSupported):
|
||||
return None
|
||||
return definition_signature(cloud_definition(wf), mapping.schedule.model_dump())
|
||||
return definition_signature(cloud_definition(wf), wire(mapping.schedule))
|
||||
|
||||
|
||||
@typechecked
|
||||
@@ -91,6 +91,11 @@ def p_mirror_cloud_state(wf: Workflow, hosted: Optional[cloud.HostedWorkflow]) -
|
||||
if hosted and wf.schedule.enabled != hosted.enabled:
|
||||
wf.schedule.enabled = hosted.enabled
|
||||
changed = True
|
||||
# Runs the cloud performed count against a max_runs cap the same as ours do, so copy its total
|
||||
# back. Without this, taking a nearly-spent schedule off the cloud would restart it at zero.
|
||||
if hosted and hosted.runs_done is not None and wf.schedule.runs_count != hosted.runs_done:
|
||||
wf.schedule.runs_count = hosted.runs_done
|
||||
changed = True
|
||||
if changed:
|
||||
storage.save_workflow(wf)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user