diff --git a/CHANGELOG.md b/CHANGELOG.md index 7f3647c..95c0f25 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,11 +7,51 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Security + +- Fixed a quoting bug in `ox_import_beat_schedules` output. Versions + 1.2.0-1.4.0 did not safely quote some stored text in printed code. + Text in the beat table could become Python that runs when the + output is applied. Only projects that ran those versions of the + importer, applied its output and had relevant beat records someone + could edit are affected, for example through Django admin beat + permissions. Installing django-ox or running the command without + applying its output does not trigger the issue. Resulting Python + could remain in `settings.py` and run whenever settings load, with + each loading process's privileges, or run with the privileges of + the process applying generated schedule calls. Version 1.5.0 + quotes every stored value with `ascii()`; tests cover every + location across SQLite, PostgreSQL and MySQL. Upgrade before + generating output. Regenerate saved output from affected versions + and review it before applying. If old output was applied, inspect + pasted settings code, deployed and repository copies, retained + schedule-call output and created schedules for unexpected Python, + tasks, arguments or timing. If unexpected Python is found or its + execution is suspected, investigate it as a security incident. + These checks cannot rule out prior execution; upgrading or + removing unexpected code does not undo it. Found during the + project's own review; there are no reports of this issue being used. + ### Fixed - Documented the `worker_class` structured log key on `claim_filter_sql_missing` and added a source-to-documentation test for structured-log extra keys. +- `ox_import_beat_schedules` now lists one-off, expired and + empty-window rows under "Not translated, and why:" instead of + importing them as recurring or live schedules. Supported schedules + retain start times only when they are later than the import instant, + otherwise starting when created; exclusive expiry bounds are always + preserved by setting `end_time` one microsecond earlier. If you + applied output from versions 1.2.0-1.4.0, review imported schedules + for unintended recurrence, missing start times and missing expiries; + disable or correct affected schedules. Database read errors, + non-null date bounds decoded as `None`, and date-bound conversion + failures stop the import with a concise error before any code is + printed. Rows with invalid JSON arguments or non-finite numeric + values are skipped with a reason. The application notes describe + naive-local date interpretation and the difference in first-run + behavior for tasks with a start time. ### Added diff --git a/docs/llms-full.txt b/docs/llms-full.txt index 5b2fffa..d6dc7a6 100644 --- a/docs/llms-full.txt +++ b/docs/llms-full.txt @@ -5025,11 +5025,51 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Security + +- Fixed a quoting bug in `ox_import_beat_schedules` output. Versions + 1.2.0-1.4.0 did not safely quote some stored text in printed code. + Text in the beat table could become Python that runs when the + output is applied. Only projects that ran those versions of the + importer, applied its output and had relevant beat records someone + could edit are affected, for example through Django admin beat + permissions. Installing django-ox or running the command without + applying its output does not trigger the issue. Resulting Python + could remain in `settings.py` and run whenever settings load, with + each loading process's privileges, or run with the privileges of + the process applying generated schedule calls. Version 1.5.0 + quotes every stored value with `ascii()`; tests cover every + location across SQLite, PostgreSQL and MySQL. Upgrade before + generating output. Regenerate saved output from affected versions + and review it before applying. If old output was applied, inspect + pasted settings code, deployed and repository copies, retained + schedule-call output and created schedules for unexpected Python, + tasks, arguments or timing. If unexpected Python is found or its + execution is suspected, investigate it as a security incident. + These checks cannot rule out prior execution; upgrading or + removing unexpected code does not undo it. Found during the + project's own review; there are no reports of this issue being used. + ### Fixed - Documented the `worker_class` structured log key on `claim_filter_sql_missing` and added a source-to-documentation test for structured-log extra keys. +- `ox_import_beat_schedules` now lists one-off, expired and + empty-window rows under "Not translated, and why:" instead of + importing them as recurring or live schedules. Supported schedules + retain start times only when they are later than the import instant, + otherwise starting when created; exclusive expiry bounds are always + preserved by setting `end_time` one microsecond earlier. If you + applied output from versions 1.2.0-1.4.0, review imported schedules + for unintended recurrence, missing start times and missing expiries; + disable or correct affected schedules. Database read errors, + non-null date bounds decoded as `None`, and date-bound conversion + failures stop the import with a concise error before any code is + printed. Rows with invalid JSON arguments or non-finite numeric + values are skipped with a reason. The application notes describe + naive-local date interpretation and the difference in first-run + behavior for tasks with a start time. ### Added diff --git a/src/django_ox/management/commands/ox_import_beat_schedules.py b/src/django_ox/management/commands/ox_import_beat_schedules.py index a27865f..2999f8c 100644 --- a/src/django_ox/management/commands/ox_import_beat_schedules.py +++ b/src/django_ox/management/commands/ox_import_beat_schedules.py @@ -9,12 +9,14 @@ from __future__ import annotations import json -from datetime import datetime +import math +from datetime import UTC, datetime, timedelta, tzinfo from typing import Any from django.conf import settings from django.core.management.base import CommandError from django.db import DatabaseError, connections, router +from django.utils import timezone from ...models import OxSchedule from .._database import DatabaseCommand @@ -28,6 +30,38 @@ CRONTAB_TABLE = "django_celery_beat_crontabschedule" INTERVAL_TABLE = "django_celery_beat_intervalschedule" +FOOTER = ( + "# Read before applying. Intervals are counted from a fixed instant here,\n" + "# not from the last run, so their fire times will differ from Celery's.\n" + "# Schedules are created with their enabled state preserved. Unless a\n" + "# start time is supplied, they start when created. If a printed start\n" + "# time passes before applying, the latest missed tick may run immediately.\n" + "# Enabling a disabled schedule resets its start time; check future starts\n" + "# before enabling.\n" + "# Start and expiry bounds are imported, with beat's exclusive expiry\n" + "# represented by an inclusive end_time one microsecond earlier. With\n" + "# USE_TZ=True, date bounds preserve their instants. With USE_TZ=False, date\n" + "# bounds are emitted as naive local times, not reinterpreted as UTC. This\n" + "# differs from django-celery-beat 2.9.0 with\n" + "# DJANGO_CELERY_BEAT_TZ_AWARE=False, which compares stored naive bounds\n" + "# against UTC wall time; review these bounds before applying the output.\n" + "# For a task that has never run and has a start_time, django-celery-beat\n" + "# 2.9.0 can run it once as soon as that time is reached or first observed\n" + "# after it has passed, then follow its schedule. The imported schedule does\n" + "# not reproduce that initial catch-up run; it waits for its first scheduled\n" + "# tick at or after start_time.\n" + "# If an expiry passes before applying, creation may fail or leave a\n" + "# schedule that never runs.\n" + "# Regenerate stale output before applying. After a partial application,\n" + "# inspect existing schedules and apply only the remaining calls; do not\n" + "# paste the whole output again.\n" + "# A queue or priority set on a beat task has no equivalent on a stored\n" + "# schedule; set it on the task." +) + +#: What _decode returns for stored arguments that are not JSON. +_INVALID = object() + PERIOD_SECONDS = { "days": 86400, "hours": 3600, @@ -68,53 +102,101 @@ def handle(self, *args: Any, **options: Any) -> None: "at the one holding your django-celery-beat schedules." ) - rows = self._read(connection) + # Every row is read and converted before a line is printed. A stored + # value the driver cannot convert, such as PostgreSQL's 'infinity', + # an impossible date kept as text on SQLite or text SQLite cannot + # decode, fails while the cursor is iterated, where no single row can + # be set aside. A date bound that cannot be carried over as it is + # stops the import the same way. So it stops, in one line and before + # any output, rather than printing a traceback or half of what it + # would have printed. + try: + rows = self._read(connection) + except (DatabaseError, ValueError, OverflowError) as exc: + raise CommandError( + f"Cannot read beat schedules from database {alias!a}; check " + "database access and stored values." + ) from exc if not rows: self.stdout.write("No periodic tasks found.") return + # One reading of the clock for the whole import, so whether a row + # is translated, the reason it is not, and the bounds its call + # carries are all decided against the same instant. + now = timezone.now() + calls = [] + skipped = [] + for row in rows: + reason = self._skip_reason(row, now) + if reason is None: + calls.append(self._as_call(row, now)) + else: + skipped.append((row["name"], reason)) + paths = sorted({row["task"] for row in rows}) self.stdout.write( "# 1. Expose these tasks. A row can only name a key you list." ) self.stdout.write('"SCHEDULABLE_TASKS": {') + # Literals through !a, like every stored value below: this fragment + # is pasted into settings. for path in paths: - self.stdout.write(f' "{path}": "{path}",') + self.stdout.write(f" {path!a}: {path!a},") self.stdout.write("},") self.stdout.write("") self.stdout.write("# 2. Create the schedules.") + self.stdout.write("from datetime import datetime") self.stdout.write("from django_ox.stored import create_schedule") self.stdout.write("") - skipped = [] - for row in rows: - line = self._as_call(row) - if line is None: - skipped.append(row) - continue - self.stdout.write(line) + for call in calls: + self.stdout.write(call) if skipped: self.stdout.write("") self.stdout.write("# Not translated, and why:") - for row in skipped: - self.stdout.write(f"# {row['name']}: {self._why(row)}") + for name, reason in skipped: + # A literal even inside a comment: a line break in a stored + # name would end the comment and paste the rest as code. + self.stdout.write(f"# {name!a}: {reason}") self.stdout.write("") - self.stdout.write( - "# Read before applying. Intervals are counted from a fixed instant " - "here, not\n# from the last run, so their fire times will differ " - "from Celery's. Schedules\n# arrive enabled and start from the " - "moment you create them. A queue or priority\n# set on a beat task " - "has no equivalent on a stored schedule; set it on the task." - ) + self.stdout.write(FOOTER) def _read(self, connection: Any) -> list[dict[str, Any]]: tables = connection.introspection.table_names() with connection.cursor() as cursor: + beat_columns = { + c.name + for c in connection.introspection.get_table_description( + cursor, BEAT_TABLE + ) + } + # expires is in django-celery-beat's first migration, one_off and + # start_time arrived in its 0007. A table older than that reads + # them as NULL rather than failing on a column it never had. + one_off_col = "one_off" if "one_off" in beat_columns else "NULL AS one_off" + start_time_col = ( + "start_time" if "start_time" in beat_columns else "NULL AS start_time" + ) + expires_col = "expires" if "expires" in beat_columns else "NULL AS expires" + # Whether each date bound is NULL, asked of the database rather + # than read off the decoded value: a driver can decode a stored + # date it cannot represent as None, as mysqlclient does with + # MySQL's zero date and Django's SQLite converter with text it + # cannot parse, and that must not read as a row without a bound. + nulls = ", ".join( + f"CASE WHEN {name} IS NULL THEN 1 ELSE 0 END AS {name}_is_null" + if name in beat_columns + else f"1 AS {name}_is_null" + for name in ("start_time", "expires") + ) + cursor.execute( f"SELECT name, task, args, kwargs, queue, enabled, " # noqa: S608 - f"crontab_id, interval_id FROM {BEAT_TABLE}" + f"crontab_id, interval_id, {one_off_col}, {start_time_col}, " + f"{expires_col}, {nulls} FROM {BEAT_TABLE}" ) columns = [c[0] for c in cursor.description] rows = [dict(zip(columns, values, strict=True)) for values in cursor] @@ -145,34 +227,123 @@ def _read(self, connection: Any) -> list[dict[str, Any]]: for row in rows: row["_cron"] = crontabs.get(row["crontab_id"]) row["_interval"] = intervals.get(row["interval_id"]) + row["_one_off"] = bool(row.get("one_off")) + row["_start_time"] = self._bound(row, "start_time", connection) + row["_expires"] = self._bound(row, "expires", connection) + row["_end_time"] = self._end_time(row["_expires"]) return rows - def _as_call(self, row: dict[str, Any]) -> str | None: - name = row["name"] + def _bound( + self, row: dict[str, Any], name: str, connection: Any + ) -> datetime | None: + """ + A row's start or expiry, which is None only where the database + holds NULL. A stored value read as no date at all is not a missing + bound, and dropping it would run the schedule outside its window. + """ + value = self._parse_datetime(row[name], connection) + if value is None and not row[f"{name}_is_null"]: + raise ValueError(f"{name} is not NULL but was read as no date") + return value + + @staticmethod + def _parse_datetime(val: Any, connection: Any) -> datetime | None: + """ + A stored start or expiry, in the form the rest of the command uses. + + With USE_TZ, Django writes a datetime to SQLite or MySQL without its + zone, in the zone of the connection: DATABASES TIME_ZONE when set, + UTC otherwise. So a naive value is read in the zone of the + connection it came from, the one --database names, not in UTC or in + TIME_ZONE. PostgreSQL returns aware values, which pass as they are. + Without USE_TZ every datetime is naive local time, and an aware one + is brought to it. + """ + if val is None: + return None + if isinstance(val, str): + val = datetime.fromisoformat(val) + if isinstance(val, datetime): + if settings.USE_TZ: + if timezone.is_naive(val): + return timezone.make_aware(val, connection.timezone) + return val + if timezone.is_aware(val): + return timezone.make_naive(val) + return val + return None + + @staticmethod + def _end_time(expires: datetime | None) -> datetime | None: + """ + The last instant a stored schedule may fire at, for a beat expiry. + + Celery stops at its expiry: a tick that falls exactly on it does not + run. A stored schedule's end_time still fires a tick that falls on + it, so the bound moves back by the smallest step a datetime holds. + The step is taken on the UTC instant rather than on the wall clock, + where a zone's clock change can skip or repeat an hour: a microsecond + before 03:00 on a day that skips from 02:00 is 01:59:59 and a + fraction, while 02:59:59 never happens and PostgreSQL would store it + an hour later. A naive expiry is local time in TIME_ZONE, so it + steps back on the instant it names there and is printed as local + time again. One that happens twice or never there names no single + instant, and is refused rather than moved. + """ + if expires is None: + return None + step = timedelta(microseconds=1) + if timezone.is_aware(expires): + return (expires.astimezone(UTC) - step).astimezone(expires.tzinfo) + zone = timezone.get_default_timezone() + if not _happens_once(expires, zone): + raise ValueError(f"{expires.isoformat()} is not one instant in {zone}") + return timezone.make_naive( + timezone.make_aware(expires, zone).astimezone(UTC) - step, zone + ) + + @staticmethod + def _future_start(row: dict[str, Any], now: datetime) -> datetime | None: + """ + The row's start time, when it is still ahead. + + create_schedule starts a schedule when it is created, which is later + than any start already past, so a past start adds nothing. + """ + start = row["_start_time"] + return start if start is not None and start > now else None + + def _as_call(self, row: dict[str, Any], now: datetime) -> str: + """The create_schedule call for a row _skip_reason lets through.""" if row["_cron"]: - minute, hour, dom, month, dow, zone = row["_cron"] - if not self._same_zone(zone): - return None - timing = f'trigger="cron", cron="{minute} {hour} {dom} {month} {dow}"' - elif row["_interval"]: + minute, hour, dom, month, dow, _zone = row["_cron"] + cron = f"{minute} {hour} {dom} {month} {dow}" + timing = f'trigger="cron", cron={cron!a}' + else: every, period = row["_interval"] seconds = every * PERIOD_SECONDS.get(period, 0) - if seconds < 1: - return None timing = f'trigger="interval", every_seconds={seconds}' - else: - return None - if self._positional_args(row): - return None arguments = self._keyword_args(row) + + start_time = self._future_start(row, now) + start_arg = "" + if start_time is not None: + start_arg = ( + f", start_time=datetime.fromisoformat({start_time.isoformat()!a})" + ) + + end_time = row["_end_time"] + end_arg = "" + if end_time is not None: + end_arg = f", end_time=datetime.fromisoformat({end_time.isoformat()!a})" + enabled = "" if row["enabled"] else ", enabled=False" - args = f", arguments={arguments!r}" if arguments else "" - # !r, not a hand-written quoted literal: a name holding a quote or a - # backslash would otherwise emit code that does not parse, or parses - # into something else, and the whole output is meant to be pasted. + args = f", arguments={arguments!a}" if arguments else "" + # Escape database-derived text with ascii() and reject non-finite numbers + # before emitting supported values as Python literals. return ( - f"create_schedule(name={name!r}, task_key={row['task']!r}, " - f"{timing}{args}{enabled})" + f"create_schedule(name={row['name']!a}, task_key={row['task']!a}, " + f"{timing}{args}{start_arg}{end_arg}{enabled})" ) @staticmethod @@ -204,10 +375,25 @@ def _same_zone(zone: Any) -> bool: @staticmethod def _decode(raw: Any) -> Any: - try: - return json.loads(raw) if isinstance(raw, str) else raw - except (TypeError, ValueError): + """ + Decode stored arguments without validating their JSON shape. + + NULL and an empty string mean no arguments, as they do to beat. + Caught JSON decoding failures return _INVALID. Valid JSON is passed + to downstream argument handling; non-object kwargs are currently + treated as empty kwargs. Stored string args are JSON-decoded and + non-string values are used unchanged; truthy results trigger the + positional-arguments reason, while falsy results, including an empty + object, are treated as no arguments. + """ + if raw is None or raw == "": return None + if not isinstance(raw, str): + return raw + try: + return json.loads(raw) + except ValueError: + return _INVALID def _positional_args(self, row: dict[str, Any]) -> bool: """A stored schedule passes keyword arguments only.""" @@ -217,25 +403,103 @@ def _keyword_args(self, row: dict[str, Any]) -> dict[str, Any]: decoded = self._decode(row.get("kwargs")) return decoded if isinstance(decoded, dict) else {} - def _why(self, row: dict[str, Any]) -> str: - if row["_cron"] and not self._same_zone(row["_cron"][5]): - return ( - f"its schedule runs in {row['_cron'][5]}, and a stored " - f"schedule has no zone of its own: it would run in " - f"{settings.TIME_ZONE}, at a different time" - ) - if row["_interval"]: + def _skip_reason(self, row: dict[str, Any], now: datetime) -> str | None: + """ + Why a row cannot be translated, or None when it can. + + The one place that decides, so a row is never printed as a call and + listed as skipped, and the reason listed is the check that failed. + """ + if row["_one_off"]: + return "one-off tasks have no equivalent on a stored schedule" + if row["_cron"]: + zone = row["_cron"][5] + if not self._same_zone(zone): + return ( + f"its schedule runs in {zone!a}, and a stored " + f"schedule has no zone of its own: it would run in " + f"{settings.TIME_ZONE!a}, at a different time" + ) + elif row["_interval"]: every, period = row["_interval"] - if every * PERIOD_SECONDS.get(period, 0) < 1: + seconds = every * PERIOD_SECONDS.get(period, 0) + if not (math.isfinite(every) and math.isfinite(seconds)): + return "interval contains a non-finite number" + if seconds < 1: return ( - f"an interval of {every} {period} is below one second, " - "which the dispatch loop cannot honour" + f"an interval of {every!a} {period!a} is below one " + "second, which the dispatch loop cannot honour" ) + elif row["crontab_id"] or row["interval_id"]: + return "its schedule row is missing" + else: + return "solar and clocked schedules have no equivalent" + for field in ("args", "kwargs"): + decoded = self._decode(row.get(field)) + if decoded is _INVALID: + return f"{field} contains invalid JSON" + if _non_finite(decoded): + return f"{field} contains a non-finite number" if self._positional_args(row): return ( "it passes positional arguments, and a stored schedule takes " "keyword arguments only; rewrite the task signature or the row" ) - if row["crontab_id"] or row["interval_id"]: - return "its schedule row is missing" - return "solar and clocked schedules have no equivalent" + return self._bounds_problem(row, now) + + def _bounds_problem(self, row: dict[str, Any], now: datetime) -> str | None: + """Why a row's start and expiry leave no window to translate, or None.""" + expires = row["_expires"] + if expires is None: + return None + # Celery counts a task as expired from the instant of its expiry. + if expires <= now: + return "it has expired" + start_time = self._future_start(row, now) + if start_time is not None and expires <= start_time: + return "its expiry is at or before its start time" + # The window the call would carry, not the one beat stored. The end + # is a microsecond before the expiry, and without a printed start the + # schedule starts when it is created, after now. create_schedule + # refuses an end that is not after the start. + if start_time is not None and row["_end_time"] <= start_time: + return ( + "its adjusted end_time is at or before its start time, too " + "short for a stored schedule" + ) + if start_time is None and row["_end_time"] <= now: + return "its expiry is one microsecond away, too short for a stored schedule" + return None + + +def _happens_once(wall: datetime, zone: tzinfo) -> bool: + """ + Does a naive local time name exactly one instant in zone? + + Where a clock change skips an hour its times never happen, and where it + repeats one they happen twice; fold picks between the two readings, and + only a time that happens once reads the same with either. + """ + return ( + wall.replace(tzinfo=zone, fold=0).utcoffset() + == wall.replace(tzinfo=zone, fold=1).utcoffset() + ) + + +def _non_finite(value: Any) -> bool: + """ + Does decoded JSON hold NaN or an infinity anywhere? json.loads accepts + both, and neither has a Python literal to be printed as. + """ + # Walk with an explicit stack: decoded JSON can be deep enough for + # a recursive non-finite check to hit Python's recursion limit. + pending = [value] + while pending: + item = pending.pop() + if isinstance(item, float) and not math.isfinite(item): + return True + if isinstance(item, dict): + pending.extend(item.values()) + elif isinstance(item, list): + pending.extend(item) + return False diff --git a/tests/test_import_beat.py b/tests/test_import_beat.py index 751b0ca..4d236b3 100644 --- a/tests/test_import_beat.py +++ b/tests/test_import_beat.py @@ -1,23 +1,47 @@ """The django-celery-beat import command, which must never write.""" +import ast +import code +import json +import sys +import tokenize +from contextlib import contextmanager, nullcontext +from datetime import UTC, datetime, timedelta, tzinfo from io import StringIO import pytest from django.conf import settings from django.core.management import call_command from django.core.management.base import CommandError -from django.db import connection +from django.db import DatabaseError, connection, connections +from django.test import override_settings +from django.utils import timezone -from django_ox.models import OxSchedule +from django_ox.models import OxSchedule, OxScheduleTick from . import tasks pytestmark = pytest.mark.django_db(transaction=True) -def make_beat_tables(): - """A minimal stand-in for the tables django-celery-beat creates.""" - with connection.cursor() as cursor: +#: The periodictask columns django-celery-beat added after its first table: +#: expires is in 0001, one_off and start_time arrived in 0007. +LATER_COLUMNS = ("one_off", "start_time", "expires") + + +def make_beat_tables(columns=LATER_COLUMNS, *, db="default"): + """ + A minimal stand-in for the tables django-celery-beat creates. + + columns picks which of the later periodictask columns exist, so an older + table is built in its own shape. The test tables are created directly + rather than reduced with DROP COLUMN, which SQLite added in 3.35.0 and + Django enables from 3.35.5. + """ + datetime_type = connections[db].data_types["DateTimeField"] + types = {"one_off": "boolean", "start_time": datetime_type} + later = "".join(f", {name} {types.get(name, datetime_type)}" for name in columns) + with connections[db].cursor() as cursor: cursor.execute( "CREATE TABLE django_celery_beat_crontabschedule (" "id integer primary key, minute varchar(64), hour varchar(64), " @@ -32,7 +56,7 @@ def make_beat_tables(): "CREATE TABLE django_celery_beat_periodictask (" "id integer primary key, name varchar(200), task varchar(200), " "args text, kwargs text, queue varchar(200), enabled boolean, " - "crontab_id integer, interval_id integer)" + f"crontab_id integer, interval_id integer{later})" ) cursor.execute( "INSERT INTO django_celery_beat_crontabschedule " @@ -45,24 +69,34 @@ def make_beat_tables(): "VALUES (1, %s, %s)", [90, "minutes"], ) - # Parameterised, and booleans passed as booleans: PostgreSQL will not - # accept 1 for a boolean column where SQLite and MySQL both would. - rows = [ - (1, "nightly", "reports.tasks.daily", None, True, 1, None), - (2, "poller", "mail.tasks.poll", "mail", True, None, 1), - (3, "orphan", "x.y.z", None, True, None, None), - ] - for pk, name, task, queue, enabled, crontab_id, interval_id in rows: - cursor.execute( - "INSERT INTO django_celery_beat_periodictask " - "(id, name, task, args, kwargs, queue, enabled, crontab_id, " - "interval_id) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s)", - [pk, name, task, "[]", "{}", queue, enabled, crontab_id, interval_id], - ) + one_off = {"one_off": False} if "one_off" in columns else {} + insert_task(1, "nightly", "reports.tasks.daily", db=db, crontab_id=1, **one_off) + insert_task( + 2, "poller", "mail.tasks.poll", db=db, queue="mail", interval_id=1, **one_off + ) + insert_task(3, "orphan", "x.y.z", db=db, **one_off) -def drop_beat_tables(): - with connection.cursor() as cursor: +def insert_task(pk, name, task, *, db="default", **columns): + """ + Add one beat row. Columns not given are NULL, or empty for the arguments. + + Parameterised, and booleans passed as booleans: PostgreSQL will not + accept 1 for a boolean column where SQLite and MySQL both would. + """ + values = {"args": "[]", "kwargs": "{}", "enabled": True, **columns} + names = ", ".join(["id", "name", "task", *values]) + marks = ", ".join(["%s"] * (len(values) + 3)) + with connections[db].cursor() as cursor: + cursor.execute( + f"INSERT INTO django_celery_beat_periodictask ({names}) " # noqa: S608 + f"VALUES ({marks})", + [pk, name, task, *values.values()], + ) + + +def drop_beat_tables(db="default"): + with connections[db].cursor() as cursor: for table in ( "django_celery_beat_periodictask", "django_celery_beat_crontabschedule", @@ -73,9 +107,13 @@ def drop_beat_tables(): @pytest.fixture def beat_tables(): - make_beat_tables() - yield - drop_beat_tables() + # Created inside the try: a setup that fails halfway still drops what + # it made, and dropping skips a table that was never created. + try: + make_beat_tables() + yield + finally: + drop_beat_tables() def run(): @@ -84,6 +122,120 @@ def run(): return out.getvalue() +SECTION_2 = "# 2. Create the schedules." + + +def section_1(output): + """The settings fragment, from its first key to section 2.""" + return output[output.index('"SCHEDULABLE_TASKS": {') : output.index(SECTION_2)] + + +def section_2(output): + """What is pasted as code: section 2's heading to the end of the output.""" + return output[output.index(SECTION_2) :] + + +def run_as_module(section): + """Section 2 as a data migration or a script runs it: compiled whole.""" + namespace = {} + exec(compile(section, "
", "exec"), namespace) # noqa: S102 + return namespace + + +class _Shell(code.InteractiveConsole): + """A manage.py shell that records what went wrong instead of printing it.""" + + def __init__(self, namespace): + super().__init__(locals=namespace) + self.errors = [] + + def showtraceback(self): + self.errors.append(sys.exc_info()[1]) + + def showsyntaxerror(self, *args, **kwargs): + self.errors.append(sys.exc_info()[1]) + + +def paste_into_shell(section): + """Section 2 as someone pastes it into a shell: one line at a time.""" + namespace = {} + shell = _Shell(namespace) + # splitlines, not split("\n"): it breaks wherever a terminal or an editor + # may end a line, at CR and U+2028 among others. + for line in section.splitlines(): + shell.push(line) + shell.push("") + assert not shell.errors, shell.errors + return namespace + + +def literal_after(text, prefix): + """The Python literal that starts right after prefix in text, evaluated.""" + rest = text[text.index(prefix) + len(prefix) :] + token = next(tokenize.generate_tokens(StringIO(rest).readline)) + return ast.literal_eval(token.string) + + +@pytest.fixture +def recorded(monkeypatch): + """create_schedule as section 2 imports it, recording its calls instead.""" + calls = [] + monkeypatch.setattr( + "django_ox.stored.create_schedule", lambda **fields: calls.append(fields) + ) + return calls + + +def calls_by_name(output, recorded): + """Section 2 run with create_schedule recorded: each call's fields by name.""" + recorded.clear() + run_as_module(section_2(output)) + return {call["name"]: call for call in recorded} + + +@pytest.fixture +def schedulable(monkeypatch): + """The fixture's two translatable task paths, registered as schedulable.""" + from django_ox.registry import ScheduleKind, register + + monkeypatch.setattr("django_ox.registry._registry", {}) + monkeypatch.setattr("django_ox.registry._discovered", True) + for key in ("reports.tasks.daily", "mail.tasks.poll"): + register(ScheduleKind(key=key, task=tasks.add)) + + +def instant(text): + """A datetime written as naive text, as the command reads it back.""" + value = datetime.fromisoformat(text) + if settings.USE_TZ: + return timezone.make_aware(value, connection.timezone) + return value + + +@pytest.fixture +def clock(monkeypatch): + """ + Pin timezone.now(), which the command and create_schedule both read. + + Returns a setter taking the time as naive text, and the list of the + instants each reading of the clock returned. + """ + readings = [] + + def pin(text): + at = instant(text) + + def now(): + readings.append(at) + return at + + monkeypatch.setattr(timezone, "now", now) + return at + + pin.readings = readings + return pin + + def test_it_writes_nothing(beat_tables): # The claim the command's own docstring makes. A migration is a decision # about production timing, so it prints and stops. @@ -93,12 +245,12 @@ def test_it_writes_nothing(beat_tables): def test_it_prints_the_allow_list_and_the_calls(beat_tables): output = run() - assert '"reports.tasks.daily": "reports.tasks.daily"' in output - assert 'cron="0 2 * * *"' in output + assert "'reports.tasks.daily': 'reports.tasks.daily'" in output + assert "cron='0 2 * * *'" in output assert "every_seconds=5400" in output -def test_the_generated_calls_actually_run(beat_tables, monkeypatch): +def test_the_generated_calls_actually_run(beat_tables, schedulable): """ Execute the output rather than matching strings in it. @@ -107,25 +259,19 @@ def test_the_generated_calls_actually_run(beat_tables, monkeypatch): pinned an output that raised TypeError the moment anyone pasted it. A printed migration is only worth printing if it runs. """ - from django_ox.registry import ScheduleKind, register - - monkeypatch.setattr("django_ox.registry._registry", {}) - monkeypatch.setattr("django_ox.registry._discovered", True) - for key in ("reports.tasks.daily", "mail.tasks.poll"): - register(ScheduleKind(key=key, task=tasks.add)) - - calls = [line for line in run().splitlines() if line.startswith("create_schedule(")] + output = run() + calls = [ + line for line in output.splitlines() if line.startswith("create_schedule(") + ] assert calls, "the command printed no calls to check" - from django_ox.stored import create_schedule - - for call in calls: - eval(call, {"create_schedule": create_schedule}) # noqa: S307 + # As pasted: on section 2's own imports, in a namespace of its own. + run_as_module(section_2(output)) assert OxSchedule.objects.count() == len(calls) -def test_a_row_with_positional_arguments_is_not_translated(beat_tables): +def test_a_row_with_positional_arguments_is_not_translated(beat_tables, recorded): # A stored schedule takes keyword arguments only, so a beat row carrying # positional args cannot be expressed and must not be printed as if it can. with connection.cursor() as cursor: @@ -134,12 +280,12 @@ def test_a_row_with_positional_arguments_is_not_translated(beat_tables): ['["emea"]'], ) output = run() - assert "nightly" not in [ - line.split("name=")[1].split(",")[0].strip("'\"") - for line in output.splitlines() - if line.startswith("create_schedule(") - ] - assert "positional arguments" in output + assert "nightly" not in calls_by_name(output, recorded) + assert ( + "# 'nightly': it passes positional arguments, and a stored schedule " + "takes keyword arguments only; rewrite the task signature or the row" + in output.splitlines() + ) def test_keyword_arguments_are_carried_over(beat_tables): @@ -152,9 +298,14 @@ def test_keyword_arguments_are_carried_over(beat_tables): def test_it_explains_what_it_could_not_translate(beat_tables): - output = run() - assert "orphan" in output - assert "Not translated" in output + # Each reason on its own row's line. These two are the fallbacks every + # other check runs before, so a check that matched too much would + # shadow them. + insert_task(4, "dangling", "x.y.z", crontab_id=99) + lines = run().splitlines() + assert "# Not translated, and why:" in lines + assert "# 'orphan': solar and clocked schedules have no equivalent" in lines + assert "# 'dangling': its schedule row is missing" in lines def test_it_warns_that_interval_timing_differs(beat_tables): @@ -163,6 +314,70 @@ def test_it_warns_that_interval_timing_differs(beat_tables): assert "fixed instant" in output +def test_it_ends_with_the_whole_note_on_applying(beat_tables): + # Keep the complete application note under regression coverage. + assert run().endswith( + "\n\n" + "# Read before applying. Intervals are counted from a fixed instant here,\n" + "# not from the last run, so their fire times will differ from Celery's.\n" + "# Schedules are created with their enabled state preserved. Unless a\n" + "# start time is supplied, they start when created. If a printed start\n" + "# time passes before applying, the latest missed tick may run immediately.\n" + "# Enabling a disabled schedule resets its start time; check future starts\n" + "# before enabling.\n" + "# Start and expiry bounds are imported, with beat's exclusive expiry\n" + "# represented by an inclusive end_time one microsecond earlier. With\n" + "# USE_TZ=True, date bounds preserve their instants. With USE_TZ=False, date\n" + "# bounds are emitted as naive local times, not reinterpreted as UTC. This\n" + "# differs from django-celery-beat 2.9.0 with\n" + "# DJANGO_CELERY_BEAT_TZ_AWARE=False, which compares stored naive bounds\n" + "# against UTC wall time; review these bounds before applying the output.\n" + "# For a task that has never run and has a start_time, django-celery-beat\n" + "# 2.9.0 can run it once as soon as that time is reached or first observed\n" + "# after it has passed, then follow its schedule. The imported schedule does\n" + "# not reproduce that initial catch-up run; it waits for its first scheduled\n" + "# tick at or after start_time.\n" + "# If an expiry passes before applying, creation may fail or leave a\n" + "# schedule that never runs.\n" + "# Regenerate stale output before applying. After a partial application,\n" + "# inspect existing schedules and apply only the remaining calls; do not\n" + "# paste the whole output again.\n" + "# A queue or priority set on a beat task has no equivalent on a stored\n" + "# schedule; set it on the task.\n" + ) + + +def note(output): + """The note after the calls, as one line of prose.""" + lines = output[output.index("# Read before applying.") :].splitlines() + return " ".join(line.removeprefix("# ") for line in lines) + + +def test_it_warns_that_beats_first_run_at_a_start_time_is_not_reproduced( + beat_tables, recorded +): + # Beat runs a task that has never run once at its start_time, off the + # schedule's ticks. The note says so, and the call carries the schedule + # as it is, with no extra run or changed cron to imitate it. + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET start_time = %s WHERE id = 1", + ["2099-01-01 10:00:00"], + ) + output = run() + calls = calls_by_name(output, recorded) + assert ( + "For a task that has never run and has a start_time, django-celery-beat " + "2.9.0 can run it once as soon as that time is reached or first observed " + "after it has passed, then follow its schedule. The imported schedule " + "does not reproduce that initial catch-up run; it waits for its first " + "scheduled tick at or after start_time." + ) in note(output) + assert sorted(call["name"] for call in recorded) == ["nightly", "poller"] + assert calls["nightly"]["cron"] == "0 2 * * *" + assert calls["nightly"]["start_time"] == instant("2099-01-01 10:00:00") + + def test_a_missing_table_is_an_error_not_an_empty_run(beat_tables): drop_beat_tables() with pytest.raises(CommandError, match="No django_celery_beat_periodictask"): @@ -170,7 +385,7 @@ def test_a_missing_table_is_an_error_not_an_empty_run(beat_tables): make_beat_tables() # so the fixture's teardown is symmetric -def test_a_schedule_in_another_timezone_is_not_translated(beat_tables): +def test_a_schedule_in_another_timezone_is_not_translated(beat_tables, recorded): # A stored schedule has no zone of its own, so a beat schedule carrying # one would run at a different time. The command names the difference # rather than emitting a line that quietly means something else. @@ -180,12 +395,12 @@ def test_a_schedule_in_another_timezone_is_not_translated(beat_tables): ["Asia/Tokyo"], ) output = run() - assert "Asia/Tokyo" in output - assert "nightly" not in [ - line.split("name=")[1].split(",")[0].strip("'\"") - for line in output.splitlines() - if line.startswith("create_schedule(") - ] + assert "nightly" not in calls_by_name(output, recorded) + assert ( + "# 'nightly': its schedule runs in 'Asia/Tokyo', and a stored schedule " + f"has no zone of its own: it would run in {settings.TIME_ZONE!a}, at a " + "different time" in output.splitlines() + ) def test_an_equivalent_zone_under_another_name_is_still_translated(beat_tables): @@ -236,3 +451,962 @@ def refuse(*args, **kwargs): with pytest.raises(CommandError) as caught: call_command("ox_import_beat_schedules") assert "Database unreachable: could not connect to server" in str(caught.value) + + +READ_ERROR = ( + "Cannot read beat schedules from database 'default'; check database access " + "and stored values." +) + +#: Stored dates some driver cannot hand back as a date, by database. +UNREADABLE = { + "postgresql": ["infinity", "-infinity"], + "sqlite": ["2099-13-45 00:00:00"], + "mysql": ["0000-00-00 00:00:00"], +} + + +def import_fails(**options): + """Run the command expecting it to stop: its one line, and what it printed.""" + out = StringIO() + with pytest.raises(CommandError) as caught: + call_command("ox_import_beat_schedules", stdout=out, **options) + return caught.value, out.getvalue() + + +def decoded_by_driver(column): + """What the driver hands back for row 1's column, or what it raises.""" + try: + with connection.cursor() as cursor: + cursor.execute( + f"SELECT {column} FROM django_celery_beat_periodictask " # noqa: S608 + "WHERE id = 1" + ) + return cursor.fetchone()[0] + except (DatabaseError, ValueError, OverflowError) as exc: + return exc + + +@pytest.mark.parametrize("column", ["start_time", "expires"]) +def test_a_date_the_driver_cannot_read_stops_the_import_in_one_line( + beat_tables, column +): + """ + PostgreSQL's 'infinity' and '-infinity', an impossible date kept as + text on SQLite, and MySQL's zero date. What reaches the command depends + on the driver: psycopg raises while the rows are read, where no single + row can be set aside; PyMySQL hands back text that is no date; and + mysqlclient hands back None for a value that is not NULL. Each stops + the import with one line, before printing anything. A driver that + decodes the value as a date, as psycopg2 does with infinity, has raised + no read error, and that value is not checked here. + """ + checked = [] + for stored in UNREADABLE[connection.vendor]: + with connection.cursor() as cursor: + cursor.execute( + f"UPDATE django_celery_beat_periodictask SET {column} = %s " # noqa: S608 + "WHERE id = 1", + [stored], + ) + if isinstance(decoded_by_driver(column), datetime): + continue + error, printed = import_fails() + assert str(error) == READ_ERROR + assert printed == "" + checked.append(stored) + if not checked: + pytest.skip( + f"this driver reads each of {UNREADABLE[connection.vendor]} as a date" + ) + + +@pytest.mark.parametrize("error", [DatabaseError, ValueError, OverflowError]) +def test_any_failure_to_read_stops_the_import_in_one_line( + beat_tables, monkeypatch, error +): + """ + The handler itself, on every backend and whatever the driver does: + database, decoding and conversion failures while the rows are read all + end in the same line, which names no driver detail. + """ + from django_ox.management.commands.ox_import_beat_schedules import Command + + def fail(*args): + raise error("a driver detail") + + monkeypatch.setattr(Command, "_parse_datetime", staticmethod(fail)) + failure, printed = import_fails() + assert str(failure) == READ_ERROR + assert type(failure.__cause__) is error + assert printed == "" + + +def test_an_expiry_in_year_one_stops_the_import_in_one_line(beat_tables, settings): + # Its end, a microsecond earlier, is before the first instant a datetime + # holds. Under UTC, so no clock change is involved. + settings.TIME_ZONE = "UTC" + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET expires = %s WHERE id = 1", + ["0001-01-01 00:00:00"], + ) + error, printed = import_fails() + assert str(error) == READ_ERROR + assert type(error.__cause__) is OverflowError + assert printed == "" + + +@pytest.mark.parametrize("column", ["start_time", "expires"]) +def test_a_stored_date_read_as_none_stops_the_import(beat_tables, column): + """ + Django's SQLite converter returns None for date text it cannot parse. + Read as no bound, a dropped start would run the schedule early and a + dropped expiry would run it forever, so the import stops instead. + """ + if connection.vendor != "sqlite": + pytest.skip("only SQLite keeps date text its converter cannot parse") + with connection.cursor() as cursor: + cursor.execute( + f"UPDATE django_celery_beat_periodictask SET {column} = %s " # noqa: S608 + "WHERE id = 1", + ["31/12/2099 10:00"], + ) + assert decoded_by_driver(column) is None + error, printed = import_fails() + assert str(error) == READ_ERROR + assert printed == "" + + +def test_a_driver_that_reads_a_stored_date_as_none_stops_the_import( + beat_tables, monkeypatch +): + # The same on every backend: a driver that decodes a stored date as None, + # as mysqlclient does with a MySQL zero date. A NULL bound still imports. + from django_ox.management.commands.ox_import_beat_schedules import Command + + monkeypatch.setattr(Command, "_parse_datetime", staticmethod(lambda *args: None)) + run() + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET expires = %s WHERE id = 1", + ["2099-12-31 10:00:00"], + ) + error, printed = import_fails() + assert str(error) == READ_ERROR + assert printed == "" + + +def test_one_off_task_is_not_translated(beat_tables, recorded): + # A stored schedule recurs, so a task meant to run once would run again. + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET one_off = %s WHERE id = 1", + [True], + ) + output = run() + assert "nightly" not in calls_by_name(output, recorded) + assert ( + "# 'nightly': one-off tasks have no equivalent on a stored schedule" + in output.splitlines() + ) + + +def test_expired_task_is_not_translated(beat_tables, recorded): + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET expires = %s WHERE id = 1", + ["2020-01-01 00:00:00"], + ) + output = run() + assert "nightly" not in calls_by_name(output, recorded) + assert "# 'nightly': it has expired" in output.splitlines() + + +def test_expiry_at_or_before_future_start_is_not_translated(beat_tables, recorded): + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET start_time = %s, " + "expires = %s WHERE id = 1", + ["2099-01-02 00:00:00", "2099-01-01 00:00:00"], + ) + output = run() + assert "nightly" not in calls_by_name(output, recorded) + assert ( + "# 'nightly': its expiry is at or before its start time" + in output.splitlines() + ) + + +def test_future_start_and_expiry_are_preserved(beat_tables, schedulable): + future_start = "2099-01-01 10:00:00" + future_expiry = "2099-12-31 10:00:00" + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET start_time = %s, " + "expires = %s WHERE id = 1", + [future_start, future_expiry], + ) + + # On section 2's own printed imports: a call that names datetime is + # only runnable if the import above it is the right one. + run_as_module(section_2(run())) + + schedule = OxSchedule.objects.get(name="nightly") + assert schedule.start_time == instant(future_start) + # A microsecond short: beat does not run a tick at its expiry, and a + # stored schedule runs one at its end_time. + assert schedule.end_time == instant(future_expiry) - timedelta(microseconds=1) + + +def test_past_start_is_omitted(beat_tables, recorded): + # create_schedule starts a schedule when it is created, which is later + # than any start already past, so a past start adds nothing. + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET start_time = %s WHERE id = 1", + ["2020-01-01 00:00:00"], + ) + assert "start_time" not in calls_by_name(run(), recorded)["nightly"] + + +@pytest.mark.parametrize( + "columns", [(), ("expires",)], ids=["none-of-them", "before-0007"] +) +def test_an_older_table_without_the_later_columns_still_imports(columns, recorded): + try: + make_beat_tables(columns) + if columns: + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET expires = %s " + "WHERE id = 1", + ["2099-12-31 10:00:00"], + ) + calls = calls_by_name(run(), recorded) + finally: + drop_beat_tables() + assert calls["nightly"]["cron"] == "0 2 * * *" + assert calls["poller"]["every_seconds"] == 5400 + assert ("end_time" in calls["nightly"]) == bool(columns) + + +def test_disabled_task_remains_disabled(beat_tables, recorded): + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET enabled = %s WHERE id = 1", + [False], + ) + calls = calls_by_name(run(), recorded) + assert calls["nightly"]["enabled"] is False + assert "enabled" not in calls["poller"] + + +def test_no_stored_value_can_leave_its_literal(beat_tables, recorded): + """ + Every stored value reaches the output as data, whatever it holds. + + The output is pasted into settings and into a shell holding production + credentials. A line break, quote or backslash in a name, a task path, a + cron field, a zone, a period or an argument must not end its literal or + its comment line. Each payload sets a marker if it ever runs as code. + """ + task = 'tasks.a"\n(OX89_TASK := "ran")\n#\\' + minute = '0" if (OX89_CRON := "ran") else "' + zone = "UTC\nOX89_ZONE = 'ran'\n#" + period = "x\u2028OX89_PERIOD = 'ran'\u2028#" + arguments = {"note": "a\rOX89_ARG = 'ran'\r#", "city": "Z\u00fcrich"} + names = { + "cron": "cr\u00f6n\nOX89_NAME = 'ran'\n#", + "interval": "interval' + (OX89_QUOTE := 'ran') + '\\", + "one-off": "one-off\u2028OX89_ONE_OFF = 'ran'\u2028#", + "expired": "expired\rOX89_EXPIRED = 'ran'\r#", + "order": "order\nOX89_ORDER = 'ran'\n#", + "orphan": 'orphan" + (OX89_ORPHAN := "ran") + "\n#', + "zone": "zone\nOX89_ZONE_NAME = 'ran'\n#", + "period": "period\rOX89_PERIOD_NAME = 'ran'\r#", + } + with connection.cursor() as cursor: + cursor.execute("DELETE FROM django_celery_beat_periodictask") + cursor.execute( + "INSERT INTO django_celery_beat_crontabschedule " + "(id, minute, hour, day_of_month, month_of_year, day_of_week, " + "timezone) VALUES (2, %s, %s, %s, %s, %s, %s), " + "(3, %s, %s, %s, %s, %s, %s)", + [minute, "2", "*", "*", "*", None, "0", "2", "*", "*", "*", zone], + ) + cursor.execute( + "INSERT INTO django_celery_beat_intervalschedule (id, every, period) " + "VALUES (2, %s, %s)", + [5, period], + ) + insert_task(1, names["cron"], task, crontab_id=2, kwargs=json.dumps(arguments)) + insert_task(2, names["interval"], task, interval_id=1) + insert_task(3, names["one-off"], task, crontab_id=1, one_off=True) + insert_task(4, names["expired"], task, interval_id=1, expires="2020-01-01 00:00:00") + insert_task( + 5, + names["order"], + task, + interval_id=1, + start_time="2099-01-02 00:00:00", + expires="2099-01-01 00:00:00", + ) + insert_task(6, names["orphan"], task) + insert_task(7, names["zone"], task, crontab_id=3) + insert_task(8, names["period"], task, interval_id=2) + + output = run() + # One physical line for every line the command wrote, whatever a + # terminal or an editor counts as a line break. + assert output.isascii() + assert "\r" not in output + + fragment = {} + exec("TASKS = {" + section_1(output) + "}", fragment) # noqa: S102 + assert fragment["TASKS"] == {"SCHEDULABLE_TASKS": {task: task}} + assert not [key for key in fragment if key.startswith("OX89")] + + for apply in (run_as_module, paste_into_shell): + recorded.clear() + namespace = apply(section_2(output)) + assert not [key for key in namespace if key.startswith("OX89")], apply + assert sorted(recorded, key=lambda call: call["name"]) == [ + { + "name": names["cron"], + "task_key": task, + "trigger": "cron", + "cron": f"{minute} 2 * * *", + "arguments": arguments, + }, + { + "name": names["interval"], + "task_key": task, + "trigger": "interval", + "every_seconds": 5400, + }, + ] + + lines = section_2(output).splitlines() + listed = lines[lines.index("# Not translated, and why:") + 1 :] + listed = [line for line in listed if line.startswith("# ")] + assert {literal_after(line, "# ") for line in listed} == { + names[key] + for key in ("one-off", "expired", "order", "orphan", "zone", "period") + } + assert literal_after(output, "its schedule runs in ") == zone + assert literal_after(output, "an interval of 5 ") == period + + +def test_one_clock_reading_decides_every_row(beat_tables, clock): + # Whether a row is translated, why it is not, and the bounds its call + # carries all come from one instant, so none can contradict another. + clock("2099-01-01 00:00:00") + insert_task(4, "expired", "x.y.z", interval_id=1, expires="2098-01-01 00:00:00") + insert_task(5, "ahead", "x.y.z", interval_id=1, start_time="2099-06-01 00:00:00") + run() + assert len(clock.readings) == 1 + + +@pytest.mark.parametrize( + ("start", "expires", "reason"), + [ + # Beat counts a task as expired from the instant of its expiry. + (None, "2099-01-01 00:00:00", "it has expired"), + ( + "2099-01-02 00:00:00", + "2099-01-02 00:00:00", + "its expiry is at or before its start time", + ), + ( + "2099-01-02 00:00:00", + "2099-01-01 23:59:59", + "its expiry is at or before its start time", + ), + # The end_time is a microsecond before the expiry, so here it would + # equal the start, and create_schedule refuses an end not after it. + ( + "2099-01-02 00:00:00", + "2099-01-02 00:00:00.000001", + ( + "its adjusted end_time is at or before its start time, too short " + "for a stored schedule" + ), + ), + # With no start printed the schedule starts when it is created, + # which is no earlier than the import. + ( + None, + "2099-01-01 00:00:00.000001", + "its expiry is one microsecond away, too short for a stored schedule", + ), + # A start at the import instant is not ahead of it, so it is left + # to create_schedule, as if there were none. + ( + "2099-01-01 00:00:00", + "2099-01-01 00:00:00.000001", + "its expiry is one microsecond away, too short for a stored schedule", + ), + ], + ids=[ + "expiry-now", + "expiry-at-start", + "expiry-before-start", + "1us", + "1us-no-start", + "1us-start-now", + ], +) +def test_a_window_no_stored_schedule_can_hold_is_listed( + beat_tables, recorded, clock, start, expires, reason +): + clock("2099-01-01 00:00:00") + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET start_time = %s, " + "expires = %s WHERE id = 1", + [start, expires], + ) + output = run() + assert "nightly" not in calls_by_name(output, recorded) + assert f"# 'nightly': {reason}" in output.splitlines() + + +@pytest.mark.parametrize( + ("start", "expires", "start_time", "end_time"), + [ + (None, "2099-01-01 00:00:00.000002", None, "2099-01-01 00:00:00.000001"), + ( + "2099-01-02 00:00:00", + "2099-01-02 00:00:00.000002", + "2099-01-02 00:00:00", + "2099-01-02 00:00:00.000001", + ), + # A start at the import instant is not ahead of it, so it is left + # out and the schedule starts when it is created. + ( + "2099-01-01 00:00:00", + "2099-01-03 00:00:00", + None, + "2099-01-02 23:59:59.999999", + ), + ], + ids=["2us-no-start", "2us-after-start", "start-now"], +) +def test_the_shortest_window_a_stored_schedule_holds_is_applied( + beat_tables, schedulable, clock, start, expires, start_time, end_time +): + now = clock("2099-01-01 00:00:00") + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET start_time = %s, " + "expires = %s WHERE id = 1", + [start, expires], + ) + output = run() + call = next(line for line in output.splitlines() if "'nightly'" in line) + assert ("start_time=" in call) == (start_time is not None) + run_as_module(section_2(output)) + schedule = OxSchedule.objects.get(name="nightly") + assert schedule.start_time == (instant(start_time) if start_time else now) + assert schedule.end_time == instant(end_time) + + +class _ClockChange(tzinfo): + """ + -05:00, then -04:00 from 2099-03-08 07:00 UTC, when local clocks skip + from 02:00 to 03:00. Built here rather than read from tz data, whose + rules for a future year can still change. + """ + + change = datetime(2099, 3, 8, 7) + + def utcoffset(self, dt): + wall = dt.replace(tzinfo=None, fold=0) + if wall < datetime(2099, 3, 8, 2): + return timedelta(hours=-5) + if wall >= datetime(2099, 3, 8, 3): + return timedelta(hours=-4) + # A wall time the change skips, read with the offset from before it + # unless fold says otherwise, as zoneinfo reads one. + return timedelta(hours=-4 if dt.fold else -5) + + def dst(self, dt): + return self.utcoffset(dt) + timedelta(hours=5) + + def tzname(self, dt): + return None + + def fromutc(self, dt): + utc = dt.replace(tzinfo=None) + offset = timedelta(hours=-5 if utc < self.change else -4) + return (utc + offset).replace(tzinfo=self) + + +def test_an_expiry_steps_back_on_its_instant_across_a_clock_change(): + # A microsecond before 03:00 on the day clocks skip an hour is 01:59:59 + # and a fraction on the wall. Stepping back on the wall clock gives + # 02:59:59, a time that never happens, read an hour after the expiry. + from django_ox.management.commands.ox_import_beat_schedules import Command + + end = Command._end_time(datetime(2099, 3, 8, 3, tzinfo=_ClockChange())) + assert end.astimezone(UTC) == datetime(2099, 3, 8, 6, 59, 59, 999999, tzinfo=UTC) + assert end.replace(tzinfo=None) == datetime(2099, 3, 8, 1, 59, 59, 999999) + + +def test_a_tick_at_the_expiry_does_not_fire(beat_tables, schedulable, clock, settings): + """ + Beat does not run a tick that falls exactly on its expiry, so the + imported schedule must not either, although it runs the tick before. + """ + from django_ox.worker import Worker + + settings.TASKS = { + "default": { + "BACKEND": "django_ox.backend.OxBackend", + "QUEUES": ["default"], + "OPTIONS": {"SCHEDULE_SOURCE": "django_ox.stored.DatabaseScheduleSource"}, + } + } + with connection.cursor() as cursor: + cursor.execute("DELETE FROM django_celery_beat_periodictask WHERE id <> 1") + # Every minute, so a tick falls exactly on an expiry on the minute. + cursor.execute( + "UPDATE django_celery_beat_crontabschedule SET minute = %s, hour = %s", + ["*", "*"], + ) + cursor.execute( + "UPDATE django_celery_beat_periodictask SET expires = %s WHERE id = 1", + ["2099-01-01 02:00:00"], + ) + clock("2099-01-01 01:00:00") + run_as_module(section_2(run())) + key = f"db:{OxSchedule.objects.get(name='nightly').pk}" + + clock("2099-01-01 01:59:00") + assert Worker(backoff_initial=0).dispatch_schedules() == 1 + expiry = clock("2099-01-01 02:00:00") + assert Worker(backoff_initial=0).dispatch_schedules() == 0 + assert list( + OxScheduleTick.objects.filter(schedule_name=key).values_list( + "scheduled_for", flat=True + ) + ) == [expiry - timedelta(minutes=1)] + + +def naive_new_york(settings): + """ + USE_TZ off in a zone with clock changes, with the worker reading stored + schedules. The dates are past ones, whose clock changes tz data will not + move, with the clock pinned before them. + """ + settings.USE_TZ = False + settings.TIME_ZONE = "America/New_York" + settings.TASKS = { + "default": { + "BACKEND": "django_ox.backend.OxBackend", + "QUEUES": ["default"], + "OPTIONS": {"SCHEDULE_SOURCE": "django_ox.stored.DatabaseScheduleSource"}, + } + } + with connection.cursor() as cursor: + cursor.execute("DELETE FROM django_celery_beat_periodictask WHERE id <> 1") + cursor.execute( + "UPDATE django_celery_beat_crontabschedule SET minute = %s, hour = %s, " + "timezone = %s", + ["*", "*", "America/New_York"], + ) + + +def test_without_use_tz_an_expiry_after_a_skipped_hour_ends_before_it( + beat_tables, schedulable, clock, settings +): + """ + At 03:00 on the day New York skips from 02:00, a microsecond before the + expiry is 01:59:59 and a fraction. On the wall clock it would be + 02:59:59, which never happens, and PostgreSQL stores that an hour later, + so the schedule would fire every minute of the hour after its expiry. + """ + from django_ox.worker import Worker + + naive_new_york(settings) + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET expires = %s WHERE id = 1", + ["2024-03-10 03:00:00"], + ) + clock("2024-03-10 00:00:00") + output = run() + assert "end_time=datetime.fromisoformat('2024-03-10T01:59:59.999999')" in output + run_as_module(section_2(output)) + schedule = OxSchedule.objects.get(name="nightly") + assert schedule.end_time == datetime(2024, 3, 10, 1, 59, 59, 999999) + fired = {} + for wall in ("01:59", "03:00", "03:01", "03:59"): + clock(f"2024-03-10 {wall}:00") + fired[wall] = Worker(backoff_initial=0).dispatch_schedules() + assert fired == {"01:59": 1, "03:00": 0, "03:01": 0, "03:59": 0} + + +@pytest.mark.parametrize( + ("expires", "end", "fired"), + [ + # 03:00 happens once, and so does the microsecond before it. + ("03:00", "02:59:59.999999", {"02:59": 1, "03:00": 0}), + # 02:00 happens once. The instant one microsecond before it is in + # the second pass through the repeated hour. The naive stored end + # is a wall-clock cutoff; these dispatch checks do not distinguish + # the two passes. + ("02:00", "01:59:59.999999", {"01:30": 1, "01:59": 1, "02:00": 0}), + ], + ids=["after-the-repeat", "end-of-the-repeat"], +) +def test_without_use_tz_an_expiry_after_a_repeated_hour_ends_before_it( + beat_tables, schedulable, clock, settings, expires, end, fired +): + """On the day New York repeats 01:00 to 02:00.""" + from django_ox.worker import Worker + + naive_new_york(settings) + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET expires = %s WHERE id = 1", + [f"2024-11-03 {expires}:00"], + ) + clock("2024-11-03 00:00:00") + output = run() + assert f"end_time=datetime.fromisoformat('2024-11-03T{end}')" in output + run_as_module(section_2(output)) + assert OxSchedule.objects.get(name="nightly").end_time == datetime.fromisoformat( + f"2024-11-03 {end}" + ) + dispatched = {} + for wall in fired: + clock(f"2024-11-03 {wall}:00") + dispatched[wall] = Worker(backoff_initial=0).dispatch_schedules() + assert dispatched == fired + + +def test_without_use_tz_a_start_in_a_skipped_hour_can_leave_no_window( + beat_tables, recorded, clock, settings +): + """ + 02:30 never happens on the day New York skips from 02:00 to 03:00, so an + expiry at 03:00 steps back to 01:59:59 and a fraction, before that start. + PostgreSQL moves the skipped start to 03:30 as it stores it, after the + expiry itself. + """ + naive_new_york(settings) + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET start_time = %s, " + "expires = %s WHERE id = 1", + ["2024-03-10 02:30:00", "2024-03-10 03:00:00"], + ) + clock("2024-03-01 00:00:00") + output = run() + assert "nightly" not in calls_by_name(output, recorded) + if connection.vendor == "postgresql": + reason = "its expiry is at or before its start time" + else: + reason = ( + "its adjusted end_time is at or before its start time, too short " + "for a stored schedule" + ) + assert f"# 'nightly': {reason}" in output.splitlines() + + +@pytest.mark.parametrize( + "expires", + ["2024-11-03 01:00:00", "2024-11-03 01:30:00", "2024-03-10 02:30:00"], + ids=["repeated-hour-start", "repeated-hour", "skipped-hour"], +) +def test_without_use_tz_an_expiry_that_is_not_one_instant_stops_the_import( + beat_tables, clock, settings, expires +): + """ + A naive local time in an hour the clocks repeat names two instants, and + one in an hour they skip names none, so no end can be derived that + means what the expiry meant. The import stops rather than guess. + """ + naive_new_york(settings) + if connection.vendor == "postgresql" and expires.startswith("2024-03-10"): + pytest.skip("PostgreSQL moves a skipped local time as it stores it") + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET expires = %s WHERE id = 1", + [expires], + ) + clock("2024-03-01 00:00:00") + error, printed = import_fails() + assert str(error) == READ_ERROR + assert printed == "" + + +@contextmanager +def connection_time_zone(alias, name): + """Point one connection at another zone, as DATABASES TIME_ZONE would.""" + wrapper = connections[alias] + original = wrapper.settings_dict["TIME_ZONE"] + + def reset(): + # Read, then delete: the zone is a cached property of the wrapper, + # and ensure_timezone sets it on the open connection, as Django's + # own suite does it. + for attr in ("timezone", "timezone_name"): + getattr(wrapper, attr) + delattr(wrapper, attr) + wrapper.ensure_timezone() + + wrapper.settings_dict["TIME_ZONE"] = name + reset() + try: + yield + finally: + wrapper.settings_dict["TIME_ZONE"] = original + reset() + + +@pytest.mark.skipif(not settings.USE_TZ, reason="a naive time has no zone") +def test_a_stored_time_is_read_in_the_connection_time_zone(beat_tables, schedulable): + """ + SQLite and MySQL keep a datetime without its zone, in the connection's + zone, and PostgreSQL returns one in it. Read in UTC instead, the + printed times would be five hours off here. A fixed offset: no clock + change, and nothing that depends on future tz data. + """ + with connection_time_zone("default", "Etc/GMT+5"): + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET start_time = %s, " + "expires = %s WHERE id = 1", + ["2099-01-01 10:00:00", "2099-12-31 10:00:00"], + ) + output = run() + # Applied in the same zone, which the stored row is read back in. + run_as_module(section_2(output)) + schedule = OxSchedule.objects.get(name="nightly") + call = next(line for line in output.splitlines() if "'nightly'" in line) + assert "start_time=datetime.fromisoformat('2099-01-01T10:00:00-05:00')" in call + assert "end_time=datetime.fromisoformat('2099-12-31T09:59:59.999999-05:00')" in call + assert schedule.start_time == datetime(2099, 1, 1, 15, tzinfo=UTC) + assert schedule.end_time == datetime(2099, 12, 31, 14, 59, 59, 999999, tzinfo=UTC) + + +@pytest.mark.skipif(not settings.USE_TZ, reason="a naive time has no zone") +@pytest.mark.django_db(transaction=True, databases=["default", "alt"]) +def test_the_zone_comes_from_the_database_it_reads(recorded): + # --database names the alias holding the beat tables, and its zone is + # the one their naive values were written in, whatever default's is. + try: + make_beat_tables(db="alt") + with connection_time_zone("alt", "Etc/GMT+5"): + with connections["alt"].cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET start_time = %s " + "WHERE id = 1", + ["2099-01-01 10:00:00"], + ) + out = StringIO() + call_command("ox_import_beat_schedules", database="alt", stdout=out) + finally: + drop_beat_tables("alt") + start = calls_by_name(out.getvalue(), recorded)["nightly"]["start_time"] + assert start == datetime(2099, 1, 1, 15, tzinfo=UTC) + assert start.utcoffset() == timedelta(hours=-5) + + +def test_without_use_tz_a_naive_local_time_is_printed_as_it_was(beat_tables, recorded): + # Without USE_TZ every datetime is naive local time, and so is the + # printed one: no offset for create_schedule to convert from. + with override_settings(USE_TZ=False) if settings.USE_TZ else nullcontext(): + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET start_time = %s, " + "expires = %s WHERE id = 1", + ["2099-01-01 10:00:00", "2099-12-31 10:00:00"], + ) + output = run() + call = next(line for line in output.splitlines() if "'nightly'" in line) + assert "start_time=datetime.fromisoformat('2099-01-01T10:00:00')" in call + assert "end_time=datetime.fromisoformat('2099-12-31T09:59:59.999999')" in call + fields = calls_by_name(output, recorded)["nightly"] + assert fields["start_time"] == datetime(2099, 1, 1, 10) + assert fields["end_time"] == datetime(2099, 12, 31, 9, 59, 59, 999999) + + +def test_without_use_tz_a_stored_offset_is_read_as_local_time(beat_tables, recorded): + # Only SQLite can hand back an offset here, from text a writer other + # than Django left. Without USE_TZ the call needs naive local time. + if connection.vendor != "sqlite": + pytest.skip("only SQLite returns an offset without USE_TZ") + with override_settings(USE_TZ=False, TIME_ZONE="Etc/GMT-3"): + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET start_time = %s " + "WHERE id = 2", + ["2099-01-01 10:00:00+00:00"], + ) + output = run() + start = calls_by_name(output, recorded)["poller"]["start_time"] + assert start == datetime(2099, 1, 1, 13) + assert start.tzinfo is None + + +@pytest.mark.skipif(settings.USE_TZ, reason="the naive settings modules") +def test_without_use_tz_a_naive_local_time_is_stored_as_it_was( + beat_tables, schedulable +): + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET start_time = %s, " + "expires = %s WHERE id = 1", + ["2099-01-01 10:00:00", "2099-12-31 10:00:00"], + ) + run_as_module(section_2(run())) + schedule = OxSchedule.objects.get(name="nightly") + assert schedule.start_time == datetime(2099, 1, 1, 10) + assert schedule.end_time == datetime(2099, 12, 31, 9, 59, 59, 999999) + + +@pytest.mark.parametrize( + ("column", "stored", "reason"), + [ + ("args", "['emea']", "args contains invalid JSON"), + ("kwargs", "{'region': 'emea'}", "kwargs contains invalid JSON"), + # More digits than Python converts from a string by default. + ("args", "[" + "7" * 5000 + "]", "args contains invalid JSON"), + ( + "kwargs", + '{"region": "emea", "n": ' + "7" * 5000 + "}", + "kwargs contains invalid JSON", + ), + ], + ids=["args-python-repr", "kwargs-python-repr", "args-long-int", "kwargs-long-int"], +) +def test_arguments_that_are_not_json_are_listed_not_dropped( + beat_tables, recorded, column, stored, reason +): + # Read as none, they would be translated into a call without them. + with connection.cursor() as cursor: + cursor.execute( + f"UPDATE django_celery_beat_periodictask SET {column} = %s " # noqa: S608 + "WHERE id = 1", + [stored], + ) + output = run() + assert "nightly" not in calls_by_name(output, recorded) + assert f"# 'nightly': {reason}" in output.splitlines() + + +@pytest.mark.parametrize("stored", [None, ""], ids=["null", "empty"]) +def test_arguments_stored_as_none_are_still_translated(beat_tables, recorded, stored): + # Beat reads NULL and an empty string as no arguments, and so does this. + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET args = %s, kwargs = %s " + "WHERE id = 1", + [stored, stored], + ) + calls = calls_by_name(run(), recorded) + assert "arguments" not in calls["nightly"] + + +@pytest.mark.parametrize( + ("column", "stored", "reason"), + [ + ("kwargs", '{"a": NaN}', "kwargs contains a non-finite number"), + ( + "kwargs", + '{"a": [1, {"b": -Infinity}]}', + "kwargs contains a non-finite number", + ), + # Finite in the text, infinite once decoded. + ("kwargs", '{"a": 1e999}', "kwargs contains a non-finite number"), + ("args", "[Infinity]", "args contains a non-finite number"), + ], + ids=["kwargs-nan", "kwargs-nested", "kwargs-overflow", "args-infinity"], +) +def test_a_non_finite_argument_is_listed_and_later_rows_still_apply( + beat_tables, recorded, column, stored, reason +): + """ + Bare nan and inf raise NameError when evaluated. In a module this + prevents later calls from running; in an interactive shell the + failing call is not applied. List the unsupported row and check + that supported rows apply in both modes. + """ + insert_task(4, "bad", "x.y.z", interval_id=1, **{column: stored}) + insert_task(5, "later", "x.y.z", interval_id=1) + output = run() + assert f"# 'bad': {reason}" in output.splitlines() + for apply in (run_as_module, paste_into_shell): + recorded.clear() + apply(section_2(output)) + assert sorted(call["name"] for call in recorded) == [ + "later", + "nightly", + "poller", + ] + + +DEEP = 900 +POSITIONAL = ( + "it passes positional arguments, and a stored schedule takes " + "keyword arguments only; rewrite the task signature or the row" +) + + +@pytest.mark.parametrize( + ("column", "stored", "reason"), + [ + ( + "args", + "[" * 600 + "0" + "]" * 600, + POSITIONAL, + ), + ( + "args", + "[" * DEEP + "0" + "]" * DEEP, + POSITIONAL, + ), + ( + "kwargs", + '{"a": ' + "[" * DEEP + "Infinity" + "]" * DEEP + "}", + "kwargs contains a non-finite number", + ), + ], + ids=["args-600", "args-900", "kwargs-900-infinity"], +) +def test_deeply_nested_arguments_are_listed_and_later_rows_still_apply( + beat_tables, recorded, column, stored, reason +): + # json.loads decodes lists nested this deep; looking through them for a + # non-finite number must not run out of recursion and stop the import. + insert_task(4, "deep", "x.y.z", interval_id=1, **{column: stored}) + insert_task(5, "later", "x.y.z", interval_id=1) + output = run() + assert f"# 'deep': {reason}" in output.splitlines() + assert sorted(calls_by_name(output, recorded)) == ["later", "nightly", "poller"] + + +@pytest.mark.parametrize( + ("every", "period"), + [(float("inf"), "minutes"), (float("-inf"), "seconds"), (1e308, "days")], + ids=["infinite", "negative-infinite", "infinite-in-seconds"], +) +def test_a_non_finite_interval_is_listed(beat_tables, recorded, every, period): + # Only SQLite can hold one, as a REAL in the integer column. + if connection.vendor != "sqlite": + pytest.skip("only SQLite stores a float in an integer column") + with connection.cursor() as cursor: + cursor.execute( + "INSERT INTO django_celery_beat_intervalschedule (id, every, period) " + "VALUES (2, %s, %s)", + [every, period], + ) + insert_task(4, "endless", "x.y.z", interval_id=2) + insert_task(5, "later", "x.y.z", interval_id=1) + output = run() + assert "# 'endless': interval contains a non-finite number" in output.splitlines() + assert sorted(calls_by_name(output, recorded)) == ["later", "nightly", "poller"]