From 5f4325dc4307d98d31343304a65836bf070ed826 Mon Sep 17 00:00:00 2001 From: Shivansh Shukla Date: Thu, 24 Sep 2026 21:11:38 +0530 Subject: [PATCH 1/2] fix: preserve supported bounds in ox_import_beat_schedules --- CHANGELOG.md | 6 + docs/llms-full.txt | 6 + .../commands/ox_import_beat_schedules.py | 86 +++++- tests/test_import_beat.py | 245 +++++++++++++++++- 4 files changed, 326 insertions(+), 17 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index f6d2a78..5d4c0cc 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,12 @@ All notable changes to this project will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [Unreleased] + +### Fixed + +- `ox_import_beat_schedules` preserves supported `start_time` and `expires` bounds when importing schedules from `django-celery-beat`. It lists one-off rows, expired rows, and expiries at or before a future start under `# Not translated, and why:`, interprets naive cursor values using `connection.timezone` when `USE_TZ=True`, and warns about expiry before application in the footer. + ## [1.4.0] - 2026-09-23 If you use Django's PostgreSQL pool, check PostgreSQL `max_connections` diff --git a/docs/llms-full.txt b/docs/llms-full.txt index d568a31..54a55da 100644 --- a/docs/llms-full.txt +++ b/docs/llms-full.txt @@ -4937,6 +4937,12 @@ All notable changes to this project will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [Unreleased] + +### Fixed + +- `ox_import_beat_schedules` preserves supported `start_time` and `expires` bounds when importing schedules from `django-celery-beat`. It lists one-off rows, expired rows, and expiries at or before a future start under `# Not translated, and why:`, interprets naive cursor values using `connection.timezone` when `USE_TZ=True`, and warns about expiry before application in the footer. + ## [1.4.0] - 2026-09-23 If you use Django's PostgreSQL pool, check PostgreSQL `max_connections` 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..60b4942 100644 --- a/src/django_ox/management/commands/ox_import_beat_schedules.py +++ b/src/django_ox/management/commands/ox_import_beat_schedules.py @@ -15,6 +15,7 @@ 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 @@ -83,6 +84,7 @@ def handle(self, *args: Any, **options: Any) -> None: 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("") @@ -104,17 +106,32 @@ def handle(self, *args: Any, **options: Any) -> None: 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." + "from Celery's. Schedules\n# preserve their enabled state and start " + "from the moment you create\n# them unless given a start time. " + "An expiry near the present can\n# pass before applying, which will " + "fail validation at creation. A queue\n# or priority set on a beat " + "task has no equivalent on a stored\n# schedule; set it on the task." ) 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 + ) + } + 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" + 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} FROM {BEAT_TABLE}" ) columns = [c[0] for c in cursor.description] rows = [dict(zip(columns, values, strict=True)) for values in cursor] @@ -145,9 +162,45 @@ 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._parse_datetime(row.get("start_time"), connection) + row["_expires"] = self._parse_datetime(row.get("expires"), connection) return rows + @staticmethod + def _parse_datetime(val: Any, connection: Any) -> datetime | None: + 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 + def _as_call(self, row: dict[str, Any]) -> str | None: + if row["_one_off"]: + return None + + now = timezone.now() + start_time = row["_start_time"] + expires = row["_expires"] + + if expires is not None and expires <= now: + return None + if ( + start_time is not None + and start_time > now + and expires is not None + and expires <= start_time + ): + return None + name = row["name"] if row["_cron"]: minute, hour, dom, month, dow, zone = row["_cron"] @@ -165,6 +218,17 @@ def _as_call(self, row: dict[str, Any]) -> str | None: if self._positional_args(row): return None arguments = self._keyword_args(row) + + start_arg = "" + if start_time is not None and start_time > now: + start_arg = ( + f", start_time=datetime.fromisoformat({start_time.isoformat()!r})" + ) + + end_arg = "" + if expires is not None and expires > now: + end_arg = f", end_time=datetime.fromisoformat({expires.isoformat()!r})" + 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 @@ -172,7 +236,7 @@ def _as_call(self, row: dict[str, Any]) -> str | None: # into something else, and the whole output is meant to be pasted. return ( f"create_schedule(name={name!r}, task_key={row['task']!r}, " - f"{timing}{args}{enabled})" + f"{timing}{args}{start_arg}{end_arg}{enabled})" ) @staticmethod @@ -218,6 +282,8 @@ def _keyword_args(self, row: dict[str, Any]) -> dict[str, Any]: return decoded if isinstance(decoded, dict) else {} def _why(self, row: dict[str, Any]) -> str: + if row["_one_off"]: + return "one-off tasks have no equivalent on a stored schedule" if row["_cron"] and not self._same_zone(row["_cron"][5]): return ( f"its schedule runs in {row['_cron'][5]}, and a stored " @@ -236,6 +302,16 @@ def _why(self, row: dict[str, Any]) -> str: "it passes positional arguments, and a stored schedule takes " "keyword arguments only; rewrite the task signature or the row" ) + now = timezone.now() + if row["_expires"] is not None and row["_expires"] <= now: + return "it has expired" + if ( + row["_start_time"] is not None + and row["_start_time"] > now + and row["_expires"] is not None + and row["_expires"] <= row["_start_time"] + ): + return "its expiry is at or before its start time" if row["crontab_id"] or row["interval_id"]: return "its schedule row is missing" return "solar and clocked schedules have no equivalent" diff --git a/tests/test_import_beat.py b/tests/test_import_beat.py index 751b0ca..2c43f1d 100644 --- a/tests/test_import_beat.py +++ b/tests/test_import_beat.py @@ -17,6 +17,7 @@ def make_beat_tables(): """A minimal stand-in for the tables django-celery-beat creates.""" + dt_type = connection.data_types.get("DateTimeField", "datetime") with connection.cursor() as cursor: cursor.execute( "CREATE TABLE django_celery_beat_crontabschedule (" @@ -29,10 +30,11 @@ def make_beat_tables(): "id integer primary key, every integer, period varchar(24))" ) cursor.execute( - "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"CREATE TABLE django_celery_beat_periodictask (" + f"id integer primary key, name varchar(200), task varchar(200), " + f"args text, kwargs text, queue varchar(200), enabled boolean, " + f"crontab_id integer, interval_id integer, one_off boolean, " + f"start_time {dt_type}, expires {dt_type})" ) cursor.execute( "INSERT INTO django_celery_beat_crontabschedule " @@ -48,16 +50,63 @@ def make_beat_tables(): # 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), + ( + 1, + "nightly", + "reports.tasks.daily", + None, + True, + 1, + None, + False, + None, + None, + ), + ( + 2, + "poller", + "mail.tasks.poll", + "mail", + True, + None, + 1, + False, + None, + None, + ), + (3, "orphan", "x.y.z", None, True, None, None, False, None, None), ] - for pk, name, task, queue, enabled, crontab_id, interval_id in rows: + for ( + pk, + name, + task, + queue, + enabled, + crontab_id, + interval_id, + one_off, + start_time, + expires, + ) 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], + "interval_id, one_off, start_time, expires) " + "VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)", + [ + pk, + name, + task, + "[]", + "{}", + queue, + enabled, + crontab_id, + interval_id, + one_off, + start_time, + expires, + ], ) @@ -107,6 +156,8 @@ 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 datetime import datetime + from django_ox.registry import ScheduleKind, register monkeypatch.setattr("django_ox.registry._registry", {}) @@ -114,13 +165,18 @@ def test_the_generated_calls_actually_run(beat_tables, monkeypatch): 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() + assert "# 2. Create the schedules." in output + section_2 = output.split("# 2. Create the schedules.")[1] + calls = [ + line for line in section_2.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 + eval(call, {"create_schedule": create_schedule, "datetime": datetime}) # noqa: S307 assert OxSchedule.objects.count() == len(calls) @@ -236,3 +292,168 @@ 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) + + +def test_one_off_task_is_not_translated(beat_tables): + 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 [ + line.split("name=")[1].split(",")[0].strip("'\"") + for line in output.splitlines() + if line.startswith("create_schedule(") + ] + assert "one-off tasks have no equivalent" in output + + +def test_expired_task_is_not_translated(beat_tables): + past_iso = "2020-01-01 00:00:00" + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET expires = %s WHERE id = 1", + [past_iso], + ) + output = run() + assert "nightly" not in [ + line.split("name=")[1].split(",")[0].strip("'\"") + for line in output.splitlines() + if line.startswith("create_schedule(") + ] + assert "it has expired" in output + + +def test_expiry_at_or_before_future_start_is_not_translated(beat_tables): + future_start = "2099-01-02 00:00:00" + earlier_expiry = "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", + [future_start, earlier_expiry], + ) + output = run() + assert "nightly" not in [ + line.split("name=")[1].split(",")[0].strip("'\"") + for line in output.splitlines() + if line.startswith("create_schedule(") + ] + assert "its expiry is at or before its start time" in output + + +def test_future_start_and_expiry_are_preserved(beat_tables, monkeypatch): + from datetime import datetime + + from django.utils import timezone + + 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)) + + 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], + ) + + output = run() + assert "start_time=datetime.fromisoformat(" in output + assert "end_time=datetime.fromisoformat(" in output + + section_2 = output.split("# 2. Create the schedules.")[1] + calls = [ + line for line in section_2.splitlines() if line.startswith("create_schedule(") + ] + + from django_ox.stored import create_schedule + + for call in calls: + eval(call, {"create_schedule": create_schedule, "datetime": datetime}) # noqa: S307 + + schedule = OxSchedule.objects.get(name="nightly") + if settings.USE_TZ: + expected_start = timezone.make_aware( + datetime.fromisoformat(future_start), connection.timezone + ) + expected_end = timezone.make_aware( + datetime.fromisoformat(future_expiry), connection.timezone + ) + else: + expected_start = datetime.fromisoformat(future_start) + expected_end = datetime.fromisoformat(future_expiry) + + assert schedule.start_time == expected_start + assert schedule.end_time == expected_end + + +def test_past_start_is_omitted(beat_tables): + past_start = "2020-01-01 00:00:00" + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET start_time = %s WHERE id = 1", + [past_start], + ) + output = run() + nightly_call = next( + line + for line in output.splitlines() + if line.startswith("create_schedule(") and "nightly" in line + ) + assert "start_time=" not in nightly_call + + +def test_old_table_without_new_columns(beat_tables): + """ + An older django-celery-beat table without one_off, start_time, or expires + must still be imported cleanly. + """ + with connection.cursor() as cursor: + cursor.execute( + "ALTER TABLE django_celery_beat_periodictask DROP COLUMN one_off" + ) + cursor.execute( + "ALTER TABLE django_celery_beat_periodictask DROP COLUMN start_time" + ) + cursor.execute( + "ALTER TABLE django_celery_beat_periodictask DROP COLUMN expires" + ) + + output = run() + assert 'cron="0 2 * * *"' in output + assert "every_seconds=5400" in output + + +def test_disabled_task_remains_disabled(beat_tables): + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET enabled = %s WHERE id = 1", + [False], + ) + output = run() + assert "enabled=False" in output + + +def test_timezone_handling_use_tz(beat_tables): + from django.test import override_settings + + future_start = "2099-05-01 12:00:00" + with connection.cursor() as cursor: + cursor.execute( + "UPDATE django_celery_beat_periodictask SET start_time = %s WHERE id = 2", + [future_start], + ) + with override_settings(USE_TZ=True, TIME_ZONE="America/New_York"): + output = run() + assert "start_time=datetime.fromisoformat(" in output + + with override_settings(USE_TZ=False, TIME_ZONE="America/New_York"): + output = run() + assert "start_time=datetime.fromisoformat(" in output From 9c8b3373ad9b469e68fcdcf9bafa04c8d0be8631 Mon Sep 17 00:00:00 2001 From: PhiLily <252857470+PhiLily@users.noreply.github.com> Date: Fri, 25 Sep 2026 12:17:59 +0300 Subject: [PATCH 2/2] Fix beat importer quoting, expiry boundaries and read failures Quote every stored value printed by ox_import_beat_schedules with ascii(), including skipped-row comments, task paths and cron fields. Keep stored text on one physical line in generated settings and calls. Set imported end_time one microsecond before beat's exclusive expiry, stepping back on the UTC instant. Without USE_TZ, resolve expiry in TIME_ZONE and print the adjusted end as local time. Stop on ambiguous or nonexistent naive expiries and on datetime underflow. Classify each row once against one import-time clock reading. List rows whose adjusted end is at or before their start, or the import instant when no start is given, with their own reason instead of a refused call. Stop before printing output when database reads, decoding or date conversion fail. Report one line naming the database, without catching programming errors. Check whether date bounds are NULL so a non-NULL bound decoded as None cannot silently lose its start or expiry. List invalid JSON arguments and non-finite argument or interval numbers with reasons. Keep NULL and empty arguments as no arguments, and retain existing handling for valid JSON with unsupported shapes. Check generated call fields using the printed imports and match reasons to individual row lines. Exercise quoted values and supported calls as a module and through InteractiveConsole. Cover older table layouts without DROP COLUMN and clean up tables after partial setup failures. Test read failures, selected connection zones, database aliases, exact stored bounds and expiry dispatch across SQLite, PostgreSQL and MySQL, including USE_TZ=False and clock changes. Correct the application note for disabled schedules, elapsed bounds, missed ticks, enabling schedules and partial application. Explain bound handling with and without USE_TZ, differences from django-celery-beat 2.9.0 with time-zone support off, and the first-run behavior not reproduced by the importer. Add a Security changelog entry for the quoting bug (CWE-94) in versions 1.2.0-1.4.0. Describe the gap in protection when output is applied, distinguish it from running the importer, and note that there are no reports of it being used. Give upgrade, regeneration and inspection guidance without implying that upgrading removes pasted code or that inspection rules out prior execution. Update the Fixed entry with untranslated rows, bound handling, read errors, application guidance and review of previously imported schedules. Regenerate docs/llms-full.txt from the changelog. Closes #86 --- CHANGELOG.md | 41 +- docs/llms-full.txt | 41 +- .../commands/ox_import_beat_schedules.py | 356 ++++- tests/test_import_beat.py | 1365 ++++++++++++++--- 4 files changed, 1511 insertions(+), 292 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 64a63f7..95c0f25 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,12 +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` preserves supported `start_time` and `expires` bounds when importing schedules from `django-celery-beat`. It lists one-off rows, expired rows, and expiries at or before a future start under `# Not translated, and why:`, interprets naive cursor values using `connection.timezone` when `USE_TZ=True`, and warns about expiry before application in the footer. +- `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 957827c..d6dc7a6 100644 --- a/docs/llms-full.txt +++ b/docs/llms-full.txt @@ -5025,12 +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` preserves supported `start_time` and `expires` bounds when importing schedules from `django-celery-beat`. It lists one-off rows, expired rows, and expiries at or before a future start under `# Not translated, and why:`, interprets naive cursor values using `connection.timezone` when `USE_TZ=True`, and warns about expiry before application in the footer. +- `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 60b4942..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,7 +9,8 @@ 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 @@ -29,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, @@ -69,18 +102,47 @@ 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.") @@ -88,30 +150,19 @@ def handle(self, *args: Any, **options: Any) -> None: 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# preserve their enabled state and start " - "from the moment you create\n# them unless given a start time. " - "An expiry near the present can\n# pass before applying, which will " - "fail validation at creation. A queue\n# or priority set on a beat " - "task has no equivalent on a stored\n# schedule; set it on the task." - ) + self.stdout.write(FOOTER) def _read(self, connection: Any) -> list[dict[str, Any]]: tables = connection.introspection.table_names() @@ -122,16 +173,30 @@ def _read(self, connection: Any) -> list[dict[str, Any]]: 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, {one_off_col}, {start_time_col}, " - f"{expires_col} FROM {BEAT_TABLE}" + 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] @@ -163,12 +228,37 @@ def _read(self, connection: Any) -> list[dict[str, Any]]: 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._parse_datetime(row.get("start_time"), connection) - row["_expires"] = self._parse_datetime(row.get("expires"), connection) + 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 _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): @@ -183,59 +273,76 @@ def _parse_datetime(val: Any, connection: Any) -> datetime | None: return val return None - def _as_call(self, row: dict[str, Any]) -> str | None: - if row["_one_off"]: + @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 + ) - now = timezone.now() - start_time = row["_start_time"] - expires = row["_expires"] + @staticmethod + def _future_start(row: dict[str, Any], now: datetime) -> datetime | None: + """ + The row's start time, when it is still ahead. - if expires is not None and expires <= now: - return None - if ( - start_time is not None - and start_time > now - and expires is not None - and expires <= start_time - ): - return None + 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 - name = row["name"] + 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 and start_time > now: + if start_time is not None: start_arg = ( - f", start_time=datetime.fromisoformat({start_time.isoformat()!r})" + f", start_time=datetime.fromisoformat({start_time.isoformat()!a})" ) + end_time = row["_end_time"] end_arg = "" - if expires is not None and expires > now: - end_arg = f", end_time=datetime.fromisoformat({expires.isoformat()!r})" + 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"create_schedule(name={row['name']!a}, task_key={row['task']!a}, " f"{timing}{args}{start_arg}{end_arg}{enabled})" ) @@ -268,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.""" @@ -281,37 +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: + 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"] 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"]: + 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" ) - now = timezone.now() - if row["_expires"] is not None and row["_expires"] <= now: + 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" - if ( - row["_start_time"] is not None - and row["_start_time"] > now - and row["_expires"] is not None - and row["_expires"] <= row["_start_time"] - ): + 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" - if row["crontab_id"] or row["interval_id"]: - return "its schedule row is missing" - return "solar and clocked schedules have no equivalent" + # 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 2c43f1d..4d236b3 100644 --- a/tests/test_import_beat.py +++ b/tests/test_import_beat.py @@ -1,24 +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.""" - dt_type = connection.data_types.get("DateTimeField", "datetime") - 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), " @@ -30,11 +53,10 @@ def make_beat_tables(): "id integer primary key, every integer, period varchar(24))" ) cursor.execute( - f"CREATE TABLE django_celery_beat_periodictask (" - f"id integer primary key, name varchar(200), task varchar(200), " - f"args text, kwargs text, queue varchar(200), enabled boolean, " - f"crontab_id integer, interval_id integer, one_off boolean, " - f"start_time {dt_type}, expires {dt_type})" + "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, " + f"crontab_id integer, interval_id integer{later})" ) cursor.execute( "INSERT INTO django_celery_beat_crontabschedule " @@ -47,71 +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, - False, - None, - None, - ), - ( - 2, - "poller", - "mail.tasks.poll", - "mail", - True, - None, - 1, - False, - None, - None, - ), - (3, "orphan", "x.y.z", None, True, None, None, False, None, None), - ] - for ( - pk, - name, - task, - queue, - enabled, - crontab_id, - interval_id, - one_off, - start_time, - expires, - ) in rows: - cursor.execute( - "INSERT INTO django_celery_beat_periodictask " - "(id, name, task, args, kwargs, queue, enabled, crontab_id, " - "interval_id, one_off, start_time, expires) " - "VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)", - [ - pk, - name, - task, - "[]", - "{}", - queue, - enabled, - crontab_id, - interval_id, - one_off, - start_time, - expires, - ], - ) + 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", @@ -122,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(): @@ -133,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. @@ -142,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. @@ -156,32 +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 datetime import datetime - - 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)) - output = run() - assert "# 2. Create the schedules." in output - section_2 = output.split("# 2. Create the schedules.")[1] calls = [ - line for line in section_2.splitlines() if line.startswith("create_schedule(") + 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, "datetime": datetime}) # 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: @@ -190,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): @@ -208,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): @@ -219,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"): @@ -226,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. @@ -236,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): @@ -294,67 +453,192 @@ def refuse(*args, **kwargs): assert "Database unreachable: could not connect to server" in str(caught.value) -def test_one_off_task_is_not_translated(beat_tables): +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 [ - line.split("name=")[1].split(",")[0].strip("'\"") - for line in output.splitlines() - if line.startswith("create_schedule(") - ] - assert "one-off tasks have no equivalent" in output + 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): - past_iso = "2020-01-01 00:00:00" +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", - [past_iso], + ["2020-01-01 00:00:00"], ) output = run() - assert "nightly" not in [ - line.split("name=")[1].split(",")[0].strip("'\"") - for line in output.splitlines() - if line.startswith("create_schedule(") - ] - assert "it has expired" in output + 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): - future_start = "2099-01-02 00:00:00" - earlier_expiry = "2099-01-01 00:00:00" +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", - [future_start, earlier_expiry], + ["2099-01-02 00:00:00", "2099-01-01 00:00:00"], ) output = run() - assert "nightly" not in [ - line.split("name=")[1].split(",")[0].strip("'\"") - for line in output.splitlines() - if line.startswith("create_schedule(") - ] - assert "its expiry is at or before its start time" in output - - -def test_future_start_and_expiry_are_preserved(beat_tables, monkeypatch): - from datetime import datetime - - from django.utils import timezone + assert "nightly" not in calls_by_name(output, recorded) + assert ( + "# 'nightly': its expiry is at or before its start time" + in output.splitlines() + ) - 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 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: @@ -364,96 +648,765 @@ def test_future_start_and_expiry_are_preserved(beat_tables, monkeypatch): [future_start, future_expiry], ) - output = run() - assert "start_time=datetime.fromisoformat(" in output - assert "end_time=datetime.fromisoformat(" in output + # 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())) - section_2 = output.split("# 2. Create the schedules.")[1] - calls = [ - line for line in section_2.splitlines() if line.startswith("create_schedule(") - ] - - from django_ox.stored import create_schedule + 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) - for call in calls: - eval(call, {"create_schedule": create_schedule, "datetime": datetime}) # noqa: S307 - schedule = OxSchedule.objects.get(name="nightly") - if settings.USE_TZ: - expected_start = timezone.make_aware( - datetime.fromisoformat(future_start), connection.timezone +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"], ) - expected_end = timezone.make_aware( - datetime.fromisoformat(future_expiry), connection.timezone + 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], ) - else: - expected_start = datetime.fromisoformat(future_start) - expected_end = datetime.fromisoformat(future_expiry) + calls = calls_by_name(run(), recorded) + assert calls["nightly"]["enabled"] is False + assert "enabled" not in calls["poller"] - assert schedule.start_time == expected_start - assert schedule.end_time == expected_end +def test_no_stored_value_can_leave_its_literal(beat_tables, recorded): + """ + Every stored value reaches the output as data, whatever it holds. -def test_past_start_is_omitted(beat_tables): - past_start = "2020-01-01 00:00:00" + 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( - "UPDATE django_celery_beat_periodictask SET start_time = %s WHERE id = 1", - [past_start], + "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], ) - output = run() - nightly_call = next( - line - for line in output.splitlines() - if line.startswith("create_schedule(") and "nightly" in line + 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", ) - assert "start_time=" not in nightly_call + 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 test_old_table_without_new_columns(beat_tables): + 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): """ - An older django-celery-beat table without one_off, start_time, or expires - must still be imported cleanly. + 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( - "ALTER TABLE django_celery_beat_periodictask DROP COLUMN one_off" + "UPDATE django_celery_beat_crontabschedule SET minute = %s, hour = %s", + ["*", "*"], ) cursor.execute( - "ALTER TABLE django_celery_beat_periodictask DROP COLUMN start_time" + "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( - "ALTER TABLE django_celery_beat_periodictask DROP COLUMN expires" + "UPDATE django_celery_beat_crontabschedule SET minute = %s, hour = %s, " + "timezone = %s", + ["*", "*", "America/New_York"], ) - output = run() - assert 'cron="0 2 * * *"' in output - assert "every_seconds=5400" in output +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 -def test_disabled_task_remains_disabled(beat_tables): + naive_new_york(settings) with connection.cursor() as cursor: cursor.execute( - "UPDATE django_celery_beat_periodictask SET enabled = %s WHERE id = 1", - [False], + "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 "enabled=False" in output + 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_timezone_handling_use_tz(beat_tables): - from django.test import override_settings +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() + - future_start = "2099-05-01 12:00:00" +@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 start_time = %s WHERE id = 2", - [future_start], + "UPDATE django_celery_beat_periodictask SET expires = %s WHERE id = 1", + [expires], ) - with override_settings(USE_TZ=True, TIME_ZONE="America/New_York"): + 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() - assert "start_time=datetime.fromisoformat(" in output - - with override_settings(USE_TZ=False, TIME_ZONE="America/New_York"): + # 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() - assert "start_time=datetime.fromisoformat(" in output + 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"]