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:]}",
|
Add the trust and safety subsystem and the App Store compliance work
Implements the moderation and consent obligations a social platform carries,
so the web version and any client that speaks to it enforce the same rules.
Moderation core (services/moderation/, database/moderation.py): a reportable
target registry, the content filter and its choke points, the report queue with
atomic resolution, enforcement actions, consent tracking, maturity gating, and
account deletion with a grace window.
Surfaces: POST /reports plus the member report list, /admin/moderation and the
per-report admin view, /workspaces, terms acceptance at /auth/terms, consent and
account deletion under /profile, the report button and dialog partials, the
maturity gate, and the moderation stylesheet and ReportDialog client.
Every user-generated surface stays reportable by construction: new content tables
are registered in REPORTABLE_TARGETS or listed in UNREPORTABLE_TABLES with a
reason, and the registry test fails the suite on anything left unclassified.
Docs: community guidelines, content moderation, intellectual property, privacy,
terms, contact, and the admin-only moderation operations page, plus the
moderation API group and the Devii moderation actions.
Compliance record: applecomp.md is the requirement register, applechanges.md the
gap analysis against this codebase, and appleimpl.md the implementation design
they resolve to.
Tests cover the report flow, admin moderation, consent, account deletion, terms
acceptance, workspaces, and the registry invariant across the unit, api, and e2e
tiers.
2026-08-09 00:18:20 +02:00
|
|
|
"terms_version": "1",
|
2026-07-26 14:57:18 +02:00
|
|
|
"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"]
|