Files
devplacepy/tests/unit/services/devii/tasks/scheduler.py
T
retoor 8e9d3fad98 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

361 lines
11 KiB
Python

# retoor <retoor@molodetz.nl>
import asyncio
from datetime import timedelta
from devplacepy.database import invalidate_admins_cache
from devplacepy.services.devii.tasks import limits
from devplacepy.services.devii.tasks.context import in_task_run
from devplacepy.services.devii.tasks.guards import (
REASON_BUDGET,
REASON_EXPIRED,
REASON_MAX_RUNS,
REASON_NOT_A_USER,
REASON_RUN_QUOTA,
)
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
SIGNUP_AFTER_PRIMARY_ADMIN = "2099-01-01T00:00:00"
def _account(local_db, role):
uid = generate_uid()
local_db["users"].insert(
{
"uid": uid,
"username": f"sched-{uid[-10:]}",
"terms_version": "1",
"role": role,
"api_key": "k",
"deleted_at": None,
"created_at": SIGNUP_AFTER_PRIMARY_ADMIN,
}
)
invalidate_admins_cache()
return uid
def _seed(local_db, owner, owner_kind="user", **overrides):
row = {
"uid": generate_uid(),
"owner_kind": owner_kind,
"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"]
def _harness(local_db, started, gate, peak=None, flags=None):
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)
if flags is not None:
flags.append(in_task_run())
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] = []
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())
assert held == 1
assert max(peak) == 1
assert len(claimed) == 1
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"
def test_scheduler_runs_a_member_task(local_db):
local_db["devii_tasks"].delete()
local_db[limits.RUNS_TABLE].delete()
owner = _account(local_db, "Member")
uid = _seed(local_db, owner)
started: list[str] = []
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())
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()
local_db[limits.RUNS_TABLE].delete()
owner = _account(local_db, "Member")
_seed(local_db, owner)
started: list[str] = []
flags: list[bool] = []
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())
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] = []
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())
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)
)
started: list[str] = []
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())
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] = []
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())
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] = []
async def run():
gate = asyncio.Event()
gate.set()
scheduler = GlobalScheduler(
local_db,
_harness(local_db, started, gate),
tick_seconds=0.05,
budget_exceeded=lambda kind, owner_id: True,
)
scheduler.start()
await asyncio.sleep(0.3)
await scheduler.stop()
run_async(run())
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] = []
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())
row = local_db["devii_tasks"].find_one(uid=uid)
assert not row["enabled"]
assert row["last_error"] == REASON_NOT_A_USER
assert started == []
def test_scheduler_retires_expired_and_exhausted_tasks_even_while_saturated(local_db):
local_db["devii_tasks"].delete()
local_db[limits.RUNS_TABLE].delete()
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] = []
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())
assert local_db["devii_tasks"].find_one(uid=expired)["last_error"] == REASON_EXPIRED
assert local_db["devii_tasks"].find_one(uid=exhausted)["last_error"] == REASON_MAX_RUNS
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"]