2026-07-26 14:57:18 +02:00
|
|
|
# retoor <retoor@molodetz.nl>
|
|
|
|
|
|
|
|
|
|
import asyncio
|
2026-07-26 16:45:38 +02:00
|
|
|
from datetime import timedelta
|
2026-07-26 14:57:18 +02:00
|
|
|
|
|
|
|
|
from devplacepy.database import invalidate_admins_cache
|
2026-07-26 16:45:38 +02:00
|
|
|
from devplacepy.services.devii.tasks import limits
|
|
|
|
|
from devplacepy.services.devii.tasks.context import in_task_run
|
2026-07-26 14:57:18 +02:00
|
|
|
from devplacepy.services.devii.tasks.guards import (
|
2026-07-26 16:45:38 +02:00
|
|
|
REASON_BUDGET,
|
2026-07-26 14:57:18 +02:00
|
|
|
REASON_EXPIRED,
|
|
|
|
|
REASON_MAX_RUNS,
|
2026-07-26 16:45:38 +02:00
|
|
|
REASON_NOT_A_USER,
|
|
|
|
|
REASON_RUN_QUOTA,
|
2026-07-26 14:57:18 +02:00
|
|
|
)
|
|
|
|
|
from devplacepy.services.devii.tasks.schedule import now_utc, to_iso
|
|
|
|
|
from devplacepy.services.devii.tasks.scheduler import GlobalScheduler
|
|
|
|
|
from devplacepy.services.devii.tasks.store import TaskStore
|
|
|
|
|
from devplacepy.utils import generate_uid
|
|
|
|
|
from tests.conftest import run_async
|
|
|
|
|
|
2026-07-26 19:58:42 +02:00
|
|
|
SIGNUP_AFTER_PRIMARY_ADMIN = "2099-01-01T00:00:00"
|
|
|
|
|
|
2026-07-26 14:57:18 +02:00
|
|
|
|
|
|
|
|
def _account(local_db, role):
|
|
|
|
|
uid = generate_uid()
|
|
|
|
|
local_db["users"].insert(
|
|
|
|
|
{
|
|
|
|
|
"uid": uid,
|
|
|
|
|
"username": f"sched-{uid[-10:]}",
|
|
|
|
|
"role": role,
|
|
|
|
|
"api_key": "k",
|
|
|
|
|
"deleted_at": None,
|
2026-07-26 19:58:42 +02:00
|
|
|
"created_at": SIGNUP_AFTER_PRIMARY_ADMIN,
|
2026-07-26 14:57:18 +02:00
|
|
|
}
|
|
|
|
|
)
|
|
|
|
|
invalidate_admins_cache()
|
|
|
|
|
return uid
|
|
|
|
|
|
|
|
|
|
|
2026-07-26 16:45:38 +02:00
|
|
|
def _seed(local_db, owner, owner_kind="user", **overrides):
|
2026-07-26 14:57:18 +02:00
|
|
|
row = {
|
|
|
|
|
"uid": generate_uid(),
|
2026-07-26 16:45:38 +02:00
|
|
|
"owner_kind": owner_kind,
|
2026-07-26 14:57:18 +02:00
|
|
|
"owner_id": owner,
|
|
|
|
|
"label": "scheduled",
|
|
|
|
|
"prompt": "work",
|
|
|
|
|
"enabled": True,
|
|
|
|
|
"status": "pending",
|
|
|
|
|
"kind": "interval",
|
|
|
|
|
"every_seconds": 900,
|
|
|
|
|
"max_runs": 10,
|
|
|
|
|
"run_count": 0,
|
|
|
|
|
"failure_count": 0,
|
|
|
|
|
"created_at": to_iso(now_utc()),
|
|
|
|
|
"next_run_at": to_iso(now_utc()),
|
|
|
|
|
"expires_at": None,
|
|
|
|
|
"deleted_at": None,
|
|
|
|
|
"deleted_by": None,
|
|
|
|
|
}
|
|
|
|
|
row.update(overrides)
|
|
|
|
|
local_db["devii_tasks"].insert(row)
|
|
|
|
|
return row["uid"]
|
|
|
|
|
|
|
|
|
|
|
2026-07-26 16:45:38 +02:00
|
|
|
def _harness(local_db, started, gate, peak=None, flags=None):
|
2026-07-26 14:57:18 +02:00
|
|
|
active: dict[str, int] = {}
|
|
|
|
|
|
|
|
|
|
def resolve(row):
|
|
|
|
|
owner = str(row["owner_id"])
|
|
|
|
|
store = TaskStore(local_db, "user", owner)
|
|
|
|
|
|
|
|
|
|
async def executor(prompt):
|
|
|
|
|
started.append(owner)
|
2026-07-26 16:45:38 +02:00
|
|
|
if flags is not None:
|
|
|
|
|
flags.append(in_task_run())
|
2026-07-26 14:57:18 +02:00
|
|
|
active[owner] = active.get(owner, 0) + 1
|
|
|
|
|
if peak is not None:
|
|
|
|
|
peak.append(max(active.values()))
|
|
|
|
|
try:
|
|
|
|
|
await gate.wait()
|
|
|
|
|
return f"done: {prompt}"
|
|
|
|
|
finally:
|
|
|
|
|
active[owner] -= 1
|
|
|
|
|
|
|
|
|
|
return store, executor, lambda kind, task_row, payload: None
|
|
|
|
|
|
|
|
|
|
return resolve
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_scheduler_runs_one_task_per_owner_and_completes_it(local_db):
|
|
|
|
|
local_db["devii_tasks"].delete()
|
|
|
|
|
owner = _account(local_db, "Admin")
|
|
|
|
|
first = _seed(local_db, owner)
|
|
|
|
|
second = _seed(local_db, owner)
|
|
|
|
|
|
|
|
|
|
started: list[str] = []
|
|
|
|
|
peak: list[int] = []
|
|
|
|
|
|
2026-07-26 19:58:42 +02:00
|
|
|
async def run():
|
|
|
|
|
gate = asyncio.Event()
|
|
|
|
|
scheduler = GlobalScheduler(
|
|
|
|
|
local_db, _harness(local_db, started, gate, peak), tick_seconds=0.05
|
|
|
|
|
)
|
|
|
|
|
scheduler.start()
|
|
|
|
|
await asyncio.sleep(0.4)
|
|
|
|
|
claimed = [r["uid"] for r in local_db["devii_tasks"].find(status="running")]
|
|
|
|
|
held = len(started)
|
|
|
|
|
gate.set()
|
|
|
|
|
await asyncio.sleep(0.5)
|
|
|
|
|
await scheduler.stop()
|
|
|
|
|
return claimed, held
|
|
|
|
|
|
|
|
|
|
claimed, held = run_async(run())
|
2026-07-26 14:57:18 +02:00
|
|
|
assert held == 1
|
|
|
|
|
assert max(peak) == 1
|
2026-07-26 19:58:42 +02:00
|
|
|
assert len(claimed) == 1
|
2026-07-26 14:57:18 +02:00
|
|
|
done = local_db["devii_tasks"].find_one(uid=claimed[0])
|
|
|
|
|
assert int(done["run_count"]) == 1
|
|
|
|
|
assert done["status"] == "pending"
|
|
|
|
|
assert done["last_result"].startswith("done:")
|
|
|
|
|
assert done["next_run_at"] > to_iso(now_utc())
|
|
|
|
|
other = second if claimed[0] == first else first
|
|
|
|
|
assert local_db["devii_tasks"].find_one(uid=other)["status"] == "pending"
|
|
|
|
|
|
|
|
|
|
|
2026-07-26 16:45:38 +02:00
|
|
|
def test_scheduler_runs_a_member_task(local_db):
|
2026-07-26 14:57:18 +02:00
|
|
|
local_db["devii_tasks"].delete()
|
2026-07-26 19:58:42 +02:00
|
|
|
local_db[limits.RUNS_TABLE].delete()
|
2026-07-26 16:45:38 +02:00
|
|
|
owner = _account(local_db, "Member")
|
|
|
|
|
uid = _seed(local_db, owner)
|
2026-07-26 14:57:18 +02:00
|
|
|
started: list[str] = []
|
|
|
|
|
|
2026-07-26 19:58:42 +02:00
|
|
|
async def run():
|
|
|
|
|
gate = asyncio.Event()
|
|
|
|
|
scheduler = GlobalScheduler(
|
|
|
|
|
local_db, _harness(local_db, started, gate), tick_seconds=0.05
|
|
|
|
|
)
|
|
|
|
|
scheduler.start()
|
|
|
|
|
await asyncio.sleep(0.3)
|
|
|
|
|
gate.set()
|
|
|
|
|
await asyncio.sleep(0.4)
|
|
|
|
|
await scheduler.stop()
|
|
|
|
|
|
|
|
|
|
run_async(run())
|
2026-07-26 16:45:38 +02:00
|
|
|
assert started == [owner]
|
|
|
|
|
row = local_db["devii_tasks"].find_one(uid=uid)
|
|
|
|
|
assert row["enabled"]
|
|
|
|
|
assert int(row["run_count"]) == 1
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_the_executor_runs_inside_the_task_run_context(local_db):
|
|
|
|
|
local_db["devii_tasks"].delete()
|
2026-07-26 19:58:42 +02:00
|
|
|
local_db[limits.RUNS_TABLE].delete()
|
2026-07-26 16:45:38 +02:00
|
|
|
owner = _account(local_db, "Member")
|
|
|
|
|
_seed(local_db, owner)
|
|
|
|
|
started: list[str] = []
|
|
|
|
|
flags: list[bool] = []
|
|
|
|
|
|
2026-07-26 19:58:42 +02:00
|
|
|
async def run():
|
|
|
|
|
gate = asyncio.Event()
|
|
|
|
|
gate.set()
|
|
|
|
|
scheduler = GlobalScheduler(
|
|
|
|
|
local_db, _harness(local_db, started, gate, flags=flags), tick_seconds=0.05
|
|
|
|
|
)
|
|
|
|
|
scheduler.start()
|
|
|
|
|
await asyncio.sleep(0.4)
|
|
|
|
|
await scheduler.stop()
|
|
|
|
|
|
|
|
|
|
run_async(run())
|
2026-07-26 16:45:38 +02:00
|
|
|
assert flags == [True]
|
|
|
|
|
assert in_task_run() is False
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_every_run_is_recorded_against_the_owner_quota(local_db):
|
|
|
|
|
local_db["devii_tasks"].delete()
|
|
|
|
|
local_db[limits.RUNS_TABLE].delete()
|
|
|
|
|
owner = _account(local_db, "Member")
|
|
|
|
|
_seed(local_db, owner)
|
|
|
|
|
started: list[str] = []
|
|
|
|
|
|
2026-07-26 19:58:42 +02:00
|
|
|
async def run():
|
|
|
|
|
gate = asyncio.Event()
|
|
|
|
|
gate.set()
|
|
|
|
|
scheduler = GlobalScheduler(
|
|
|
|
|
local_db, _harness(local_db, started, gate), tick_seconds=0.05
|
|
|
|
|
)
|
|
|
|
|
scheduler.start()
|
|
|
|
|
await asyncio.sleep(0.4)
|
|
|
|
|
await scheduler.stop()
|
|
|
|
|
|
|
|
|
|
run_async(run())
|
2026-07-26 16:45:38 +02:00
|
|
|
quota = limits.run_quota(local_db, "user", owner, now_utc())
|
|
|
|
|
assert quota.used == 1
|
|
|
|
|
assert quota.limit == limits.DEFAULT_MEMBER_RUNS
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_a_member_over_the_run_quota_is_postponed_not_disabled(local_db):
|
|
|
|
|
local_db["devii_tasks"].delete()
|
|
|
|
|
local_db[limits.RUNS_TABLE].delete()
|
|
|
|
|
owner = _account(local_db, "Member")
|
|
|
|
|
uid = _seed(local_db, owner)
|
|
|
|
|
now = now_utc()
|
|
|
|
|
for index in range(limits.DEFAULT_MEMBER_RUNS):
|
|
|
|
|
limits.record_run(
|
|
|
|
|
local_db, "user", owner, generate_uid(), now - timedelta(minutes=index)
|
2026-07-26 14:57:18 +02:00
|
|
|
)
|
2026-07-26 16:45:38 +02:00
|
|
|
started: list[str] = []
|
2026-07-26 14:57:18 +02:00
|
|
|
|
2026-07-26 19:58:42 +02:00
|
|
|
async def run():
|
|
|
|
|
gate = asyncio.Event()
|
|
|
|
|
gate.set()
|
|
|
|
|
scheduler = GlobalScheduler(
|
|
|
|
|
local_db, _harness(local_db, started, gate), tick_seconds=0.05
|
|
|
|
|
)
|
|
|
|
|
scheduler.start()
|
|
|
|
|
await asyncio.sleep(0.4)
|
|
|
|
|
await scheduler.stop()
|
|
|
|
|
|
|
|
|
|
run_async(run())
|
2026-07-26 16:45:38 +02:00
|
|
|
row = local_db["devii_tasks"].find_one(uid=uid)
|
|
|
|
|
assert started == []
|
|
|
|
|
assert row["enabled"]
|
|
|
|
|
assert row["status"] == "pending"
|
|
|
|
|
assert row["last_error"] == REASON_RUN_QUOTA
|
|
|
|
|
assert row["next_run_at"] > to_iso(now)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_an_administrator_keeps_running_past_the_member_limit(local_db):
|
|
|
|
|
local_db["devii_tasks"].delete()
|
|
|
|
|
local_db[limits.RUNS_TABLE].delete()
|
|
|
|
|
owner = _account(local_db, "Admin")
|
|
|
|
|
_seed(local_db, owner)
|
|
|
|
|
now = now_utc()
|
|
|
|
|
for index in range(limits.DEFAULT_MEMBER_RUNS + 5):
|
|
|
|
|
limits.record_run(
|
|
|
|
|
local_db, "user", owner, generate_uid(), now - timedelta(minutes=index)
|
|
|
|
|
)
|
|
|
|
|
started: list[str] = []
|
|
|
|
|
|
2026-07-26 19:58:42 +02:00
|
|
|
async def run():
|
|
|
|
|
gate = asyncio.Event()
|
|
|
|
|
gate.set()
|
|
|
|
|
scheduler = GlobalScheduler(
|
|
|
|
|
local_db, _harness(local_db, started, gate), tick_seconds=0.05
|
|
|
|
|
)
|
|
|
|
|
scheduler.start()
|
|
|
|
|
await asyncio.sleep(0.4)
|
|
|
|
|
await scheduler.stop()
|
|
|
|
|
|
|
|
|
|
run_async(run())
|
2026-07-26 16:45:38 +02:00
|
|
|
assert started == [owner]
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_budget_postpones_the_task(local_db):
|
|
|
|
|
local_db["devii_tasks"].delete()
|
|
|
|
|
local_db[limits.RUNS_TABLE].delete()
|
|
|
|
|
owner = _account(local_db, "Admin")
|
|
|
|
|
uid = _seed(local_db, owner)
|
|
|
|
|
started: list[str] = []
|
|
|
|
|
|
2026-07-26 19:58:42 +02:00
|
|
|
async def run():
|
|
|
|
|
gate = asyncio.Event()
|
|
|
|
|
gate.set()
|
|
|
|
|
scheduler = GlobalScheduler(
|
2026-07-26 16:45:38 +02:00
|
|
|
local_db,
|
|
|
|
|
_harness(local_db, started, gate),
|
2026-07-26 19:58:42 +02:00
|
|
|
tick_seconds=0.05,
|
2026-07-26 16:45:38 +02:00
|
|
|
budget_exceeded=lambda kind, owner_id: True,
|
2026-07-26 19:58:42 +02:00
|
|
|
)
|
|
|
|
|
scheduler.start()
|
|
|
|
|
await asyncio.sleep(0.3)
|
|
|
|
|
await scheduler.stop()
|
|
|
|
|
|
|
|
|
|
run_async(run())
|
2026-07-26 16:45:38 +02:00
|
|
|
row = local_db["devii_tasks"].find_one(uid=uid)
|
|
|
|
|
assert started == []
|
|
|
|
|
assert row["enabled"]
|
|
|
|
|
assert row["last_error"] == REASON_BUDGET
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_scheduler_retires_a_guest_owned_task(local_db):
|
|
|
|
|
local_db["devii_tasks"].delete()
|
|
|
|
|
uid = _seed(local_db, "guest-cookie", owner_kind="guest")
|
|
|
|
|
started: list[str] = []
|
|
|
|
|
|
2026-07-26 19:58:42 +02:00
|
|
|
async def run():
|
|
|
|
|
gate = asyncio.Event()
|
|
|
|
|
gate.set()
|
|
|
|
|
scheduler = GlobalScheduler(
|
|
|
|
|
local_db, _harness(local_db, started, gate), tick_seconds=0.05
|
|
|
|
|
)
|
|
|
|
|
scheduler.start()
|
|
|
|
|
await asyncio.sleep(0.3)
|
|
|
|
|
await scheduler.stop()
|
|
|
|
|
|
|
|
|
|
run_async(run())
|
2026-07-26 14:57:18 +02:00
|
|
|
row = local_db["devii_tasks"].find_one(uid=uid)
|
|
|
|
|
assert not row["enabled"]
|
2026-07-26 16:45:38 +02:00
|
|
|
assert row["last_error"] == REASON_NOT_A_USER
|
2026-07-26 14:57:18 +02:00
|
|
|
assert started == []
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_scheduler_retires_expired_and_exhausted_tasks_even_while_saturated(local_db):
|
|
|
|
|
local_db["devii_tasks"].delete()
|
2026-07-26 16:45:38 +02:00
|
|
|
local_db[limits.RUNS_TABLE].delete()
|
2026-07-26 14:57:18 +02:00
|
|
|
owner = _account(local_db, "Admin")
|
|
|
|
|
_seed(local_db, owner)
|
|
|
|
|
expired = _seed(
|
|
|
|
|
local_db, owner, expires_at=to_iso(now_utc().replace(year=now_utc().year - 1))
|
|
|
|
|
)
|
|
|
|
|
exhausted = _seed(local_db, owner, max_runs=3, run_count=3)
|
|
|
|
|
started: list[str] = []
|
|
|
|
|
|
2026-07-26 19:58:42 +02:00
|
|
|
async def run():
|
|
|
|
|
gate = asyncio.Event()
|
|
|
|
|
scheduler = GlobalScheduler(
|
|
|
|
|
local_db, _harness(local_db, started, gate), tick_seconds=0.05
|
|
|
|
|
)
|
|
|
|
|
scheduler.start()
|
|
|
|
|
await asyncio.sleep(0.4)
|
|
|
|
|
gate.set()
|
|
|
|
|
await asyncio.sleep(0.2)
|
|
|
|
|
await scheduler.stop()
|
|
|
|
|
|
|
|
|
|
run_async(run())
|
2026-07-26 14:57:18 +02:00
|
|
|
assert local_db["devii_tasks"].find_one(uid=expired)["last_error"] == REASON_EXPIRED
|
2026-07-26 19:58:42 +02:00
|
|
|
assert local_db["devii_tasks"].find_one(uid=exhausted)["last_error"] == REASON_MAX_RUNS
|
2026-07-26 14:57:18 +02:00
|
|
|
for uid in (expired, exhausted):
|
|
|
|
|
assert not local_db["devii_tasks"].find_one(uid=uid)["enabled"]
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_scheduler_disables_a_task_after_repeated_failures(local_db):
|
|
|
|
|
local_db["devii_tasks"].delete()
|
|
|
|
|
owner = _account(local_db, "Admin")
|
|
|
|
|
uid = _seed(local_db, owner, failure_count=2)
|
|
|
|
|
|
|
|
|
|
def resolve(row):
|
|
|
|
|
store = TaskStore(local_db, "user", str(row["owner_id"]))
|
|
|
|
|
|
|
|
|
|
async def executor(prompt):
|
|
|
|
|
raise RuntimeError("upstream is down")
|
|
|
|
|
|
|
|
|
|
return store, executor, lambda kind, task_row, payload: None
|
|
|
|
|
|
|
|
|
|
async def run():
|
|
|
|
|
scheduler = GlobalScheduler(
|
|
|
|
|
local_db, resolve, tick_seconds=0.05, max_failures=3
|
|
|
|
|
)
|
|
|
|
|
scheduler.start()
|
|
|
|
|
await asyncio.sleep(0.4)
|
|
|
|
|
await scheduler.stop()
|
|
|
|
|
|
|
|
|
|
run_async(run())
|
|
|
|
|
row = local_db["devii_tasks"].find_one(uid=uid)
|
|
|
|
|
assert int(row["failure_count"]) == 3
|
|
|
|
|
assert not row["enabled"]
|
|
|
|
|
assert row["status"] == "error"
|
|
|
|
|
assert "upstream is down" in row["last_error"]
|