Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
103 changes: 103 additions & 0 deletions .semgrep/rules/devex/schedule-must-avoid-minute-zero.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,103 @@
# Test cases for schedule-must-avoid-minute-zero.
# ruff: noqa
import datetime as dt
import os
from datetime import timedelta

import dagster
from celery.schedules import crontab
from temporalio.client import ScheduleIntervalSpec, ScheduleSpec

from posthog.scheduling.jitter import deterministic_offset

# ruleid: schedule-must-avoid-minute-zero
sender.add_periodic_task(crontab(hour="*", minute="0"), refresh_cache.s())

# ruleid: schedule-must-avoid-minute-zero
sender.add_periodic_task(crontab(minute="0"), kill_stale_runs.s())

# ruleid: schedule-must-avoid-minute-zero
sender.add_periodic_task(crontab(minute="0", hour="*/12"), refresh_fields.s())

# ruleid: schedule-must-avoid-minute-zero
sender.add_periodic_task(crontab(minute=0), sweep.s())

# ok: schedule-must-avoid-minute-zero
sender.add_periodic_task(crontab(hour="*", minute="23"), refresh_cache.s())

# ok: schedule-must-avoid-minute-zero
sender.add_periodic_task(crontab(hour="3", minute="0"), daily_cleanup.s())

# ok: schedule-must-avoid-minute-zero
sender.add_periodic_task(crontab(minute="*/5"), poll.s())


# ruleid: schedule-must-avoid-minute-zero
@dagster.schedule(cron_schedule="0 * * * *", job=hourly_job)
def hourly_schedule(context):
return dagster.RunRequest()


# ruleid: schedule-must-avoid-minute-zero
every_six_hours = dagster.ScheduleDefinition(job=job, cron_schedule="0 */6 * * *")


# ok: schedule-must-avoid-minute-zero
@dagster.schedule(cron_schedule="17 * * * *", job=hourly_job)
def offset_schedule(context):
return dagster.RunRequest()


# ok: schedule-must-avoid-minute-zero
daily = dagster.ScheduleDefinition(job=job, cron_schedule="0 3 * * *")

# ruleid: schedule-must-avoid-minute-zero
HOURLY_CRON_SCHEDULE = os.getenv("HOURLY_CRON_SCHEDULE", "0 * * * *")

# ok: schedule-must-avoid-minute-zero
named_hourly = dagster.ScheduleDefinition(job=job, cron_schedule=HOURLY_CRON_SCHEDULE)

# ruleid: schedule-must-avoid-minute-zero
hourly_cron = ScheduleSpec(cron_expressions=["0 * * * *"])

# ok: schedule-must-avoid-minute-zero
every_minute = ScheduleSpec(cron_expressions=["*/1 * * * *"])

# ok: schedule-must-avoid-minute-zero
daily_cron = ScheduleSpec(cron_expressions=["2 3 * * *"], jitter=timedelta(minutes=30))

# ruleid: schedule-must-avoid-minute-zero
hourly_interval = ScheduleIntervalSpec(every=timedelta(hours=1))

# ruleid: schedule-must-avoid-minute-zero
quarter_hour = ScheduleIntervalSpec(every=dt.timedelta(minutes=15))

# ruleid: schedule-must-avoid-minute-zero
named_interval = ScheduleIntervalSpec(every=SCHEDULE_INTERVAL)

# ok: schedule-must-avoid-minute-zero
offset_interval = ScheduleIntervalSpec(every=timedelta(hours=1), offset=timedelta(minutes=2))

# ruleid: schedule-must-avoid-minute-zero
none_offset = ScheduleIntervalSpec(every=timedelta(hours=1), offset=None)

# ruleid: schedule-must-avoid-minute-zero
empty_offset = ScheduleIntervalSpec(every=timedelta(hours=1), offset=timedelta())

# ruleid: schedule-must-avoid-minute-zero
zero_offset = ScheduleIntervalSpec(every=timedelta(hours=1), offset=timedelta(0))

# ruleid: schedule-must-avoid-minute-zero
zero_minutes_offset = ScheduleIntervalSpec(every=dt.timedelta(hours=1), offset=dt.timedelta(minutes=0))

# ok: schedule-must-avoid-minute-zero
entity_interval = ScheduleIntervalSpec(every=SCHEDULE_INTERVAL, offset=deterministic_offset(key, SCHEDULE_INTERVAL))

# ok: schedule-must-avoid-minute-zero
daily_interval = ScheduleIntervalSpec(every=timedelta(days=1))

# ruleid: schedule-must-avoid-minute-zero
six_hourly_interval = ScheduleIntervalSpec(every=timedelta(hours=6))

# ok: schedule-must-avoid-minute-zero
day_in_hours_interval = ScheduleIntervalSpec(every=timedelta(hours=24))
98 changes: 98 additions & 0 deletions .semgrep/rules/devex/schedule-must-avoid-minute-zero.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
rules:
- id: schedule-must-avoid-minute-zero
message: |
This schedule fires at minute zero of the hour.

Fleet schedules that fire together concentrate query load. Check how the job tracks
its input before moving it: a rolling lookback equal to the interval can skip data
when its start moves. Jitter needs a cursor or enough overlapping lookback.

- Celery beat and Dagster: write a fixed minute that no other job in that hour holds,
for example `crontab(minute="23")` or `cron_schedule="23 * * * *"`.
- Temporal fleet schedule: start two minutes after the boundary and add jitter,
for example `ScheduleIntervalSpec(every=timedelta(hours=1), offset=timedelta(minutes=2))`
with `ScheduleSpec(..., jitter=timedelta(minutes=10))`.
- Sub-hourly or per-entity Temporal schedule: pass
`offset=deterministic_offset(<stable id>, <interval>)` from `posthog/scheduling/jitter.py`.

Keep the existing time when a customer chose it or a fixed lookback requires it. Add
`# nosemgrep: schedule-must-avoid-minute-zero -- <reason>`.
languages: [python]
severity: ERROR
Comment thread
aspicer marked this conversation as resolved.
pattern-either:
# Celery beat at minute 0 of every hour, or of every N hours. crontab's hour defaults to "*".
- patterns:
- pattern-either:
- pattern: crontab(..., minute="0", ...)
- pattern: crontab(..., minute="00", ...)
- pattern: crontab(..., minute=0, ...)
- pattern-either:
- patterns:
- pattern: crontab(...)
- pattern-not: crontab(..., hour=$HOUR, ...)
- patterns:
- pattern: crontab(..., hour="$HOUR", ...)
- metavariable-regex:
metavariable: $HOUR
regex: ^\*(/\d+)?$
# A five-field cron string with minute 0 and an hour field of "*" or "*/N". The rule matches the
# literal wherever it is written, because a Dagster or Temporal schedule often reads it from a
# constant or an environment default in another module, which semgrep cannot follow.
- patterns:
- pattern: '"$CRON"'
- metavariable-regex:
metavariable: $CRON
regex: ^\s*0+\s+\*(/\d+)?(\s+\S+){3}\s*$
# A Temporal interval with no offset, or a zero offset, starts on the epoch grid, so a sub-daily
# interval that fits evenly into the day fires at minute zero. An interval the rule cannot read
# counts as short.
- patterns:
- pattern-either:
- patterns:
- pattern: ScheduleIntervalSpec(...)
- pattern-not: ScheduleIntervalSpec(..., offset=$OFFSET, ...)
- pattern: ScheduleIntervalSpec(..., offset=None, ...)
- pattern: ScheduleIntervalSpec(..., offset=datetime.timedelta(), ...)
- pattern: ScheduleIntervalSpec(..., offset=datetime.timedelta(0), ...)
- pattern: ScheduleIntervalSpec(..., offset=datetime.timedelta($UNIT=0), ...)
- pattern-either:
- patterns:
- pattern: ScheduleIntervalSpec(..., every=$INTERVAL, ...)
- metavariable-regex:
metavariable: $INTERVAL
regex: ^[A-Za-z_][\w.]*$
- patterns:
- pattern: ScheduleIntervalSpec(..., every=datetime.timedelta(hours=$HOURS), ...)
- metavariable-regex:
metavariable: $HOURS
regex: ^([1-9]|1\d|2[0-3]|[A-Za-z_][\w.]*)$
- patterns:
- pattern: ScheduleIntervalSpec(..., every=datetime.timedelta(minutes=$MINUTES), ...)
- metavariable-regex:
metavariable: $MINUTES
regex: ^([1-9]\d{0,2}|1[0-3]\d\d|14[0-3]\d|[A-Za-z_][\w.]*)$
- patterns:
- pattern: ScheduleIntervalSpec(..., every=datetime.timedelta(seconds=$SECONDS), ...)
- metavariable-regex:
metavariable: $SECONDS
regex: ^([1-9]\d{0,3}|[1-7]\d{4}|8[0-5]\d{3}|86[0-3]\d\d|[A-Za-z_][\w.]*)$
# Constant propagation also matches each use of a cron constant, so one schedule would need a
# nosemgrep comment at the definition and at every use. Without it, the rule reports the definition only.
options:
constant_propagation: false
paths:
include:
- dags/**/*.py
- ee/**/*.py
- posthog/**/*.py
- products/**/*.py
exclude:
- '**/test/**'
- '**/tests/**'
- '**/test_*.py'
- '**/*_test.py'
- '**/conftest*.py'
metadata:
category: performance
subcategory: audit
confidence: HIGH
1 change: 1 addition & 0 deletions posthog/temporal/ai_observability/eval_reports/schedule.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ async def create_eval_reports_schedule(client: Client):
id=SCHEDULE_ID,
task_queue=settings.LLMA_TASK_QUEUE,
),
# nosemgrep: schedule-must-avoid-minute-zero -- delivers reports at the times customers schedule them, with a 15-minute lookahead
spec=ScheduleSpec(intervals=[ScheduleIntervalSpec(every=timedelta(hours=1))]),
)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ async def create_evaluation_sampler_schedule(client: Client) -> None:
task_queue=settings.LLMA_TASK_QUEUE,
execution_timeout=SAMPLER_COORDINATOR_EXECUTION_TIMEOUT,
),
# nosemgrep: schedule-must-avoid-minute-zero -- the rolling one-hour sample window has no cursor, so shifting the schedule skips data
spec=ScheduleSpec(intervals=[ScheduleIntervalSpec(every=timedelta(hours=SAMPLER_SCHEDULE_INTERVAL_HOURS))]),
policy=SchedulePolicy(overlap=ScheduleOverlapPolicy.SKIP),
)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@ async def create_batch_trace_summarization_schedule(client: Client):
task_queue=settings.LLMA_TASK_QUEUE,
execution_timeout=timedelta(minutes=COORDINATOR_EXECUTION_TIMEOUT_MINUTES),
),
# nosemgrep: schedule-must-avoid-minute-zero -- the rolling one-hour lookback has no cursor, so shifting the schedule skips data
spec=ScheduleSpec(intervals=[ScheduleIntervalSpec(every=timedelta(hours=SCHEDULE_INTERVAL_HOURS))]),
policy=SchedulePolicy(overlap=ScheduleOverlapPolicy.SKIP),
)
Expand Down
1 change: 1 addition & 0 deletions posthog/temporal/schedule.py
Original file line number Diff line number Diff line change
Expand Up @@ -661,6 +661,7 @@ async def create_replay_count_metrics_schedule(client: Client) -> None:
),
),
spec=ScheduleSpec(
# nosemgrep: schedule-must-avoid-minute-zero -- the metric query reads a rolling hour with no cursor, so shifting the schedule skips data
intervals=[ScheduleIntervalSpec(every=timedelta(hours=1))],
),
)
Expand Down
37 changes: 36 additions & 1 deletion posthog/temporal/tests/test_schedule.py
Original file line number Diff line number Diff line change
@@ -1,15 +1,23 @@
from collections.abc import Awaitable, Callable
from datetime import timedelta

import pytest
from unittest import mock

from django.conf import settings
from django.test import override_settings

from temporalio.client import Schedule, ScheduleActionStartWorkflow
from temporalio.client import Client, Schedule, ScheduleActionStartWorkflow
from temporalio.service import RPCError, RPCStatusCode

from posthog.temporal.ai_observability.evaluation_clustering.constants import SAMPLER_WINDOW_MINUTES
from posthog.temporal.ai_observability.trace_clustering import constants as trace_clustering_constants
from posthog.temporal.ai_observability.trace_summarization import constants as trace_summarization_constants
from posthog.temporal.schedule import (
cleanup_non_cloud_ai_observability_schedules,
create_batch_trace_summarization_schedule,
create_evaluation_sampler_schedule,
create_replay_count_metrics_schedule,
create_wa_digest_notification_schedule,
create_wa_weekly_digest_schedule,
)
Expand Down Expand Up @@ -79,3 +87,30 @@ async def test_cleanup_non_cloud_ai_observability_schedules(cloud_deployment, de
await cleanup_non_cloud_ai_observability_schedules(mock.MagicMock())

assert deleted == expected_deleted


@pytest.mark.asyncio
@pytest.mark.parametrize(
"create_schedule,lookback",
[
(
create_batch_trace_summarization_schedule,
timedelta(minutes=trace_summarization_constants.DEFAULT_WINDOW_MINUTES),
),
(create_evaluation_sampler_schedule, timedelta(minutes=SAMPLER_WINDOW_MINUTES)),
(create_replay_count_metrics_schedule, timedelta(hours=1)),
],
)
async def test_fixed_window_schedules_preserve_coverage(
create_schedule: Callable[[Client], Awaitable[None]], lookback: timedelta
) -> None:
client = mock.MagicMock(spec=Client)
client.get_schedule_handle.return_value.describe = mock.AsyncMock(
side_effect=RPCError("not found", RPCStatusCode.NOT_FOUND, b"")
)
await create_schedule(client)

schedule = client.create_schedule.call_args.kwargs["schedule"]
interval = schedule.spec.intervals[0]
assert interval.every + (schedule.spec.jitter or timedelta(0)) <= lookback
assert (interval.offset or timedelta(0)) == timedelta(0)
1 change: 1 addition & 0 deletions products/batch_exports/backend/service.py
Original file line number Diff line number Diff line change
Expand Up @@ -1168,6 +1168,7 @@ def _get_schedule_spec(batch_export: BatchExport) -> ScheduleSpec:
return ScheduleSpec(
start_at=batch_export.start_at,
end_at=batch_export.end_at,
# nosemgrep: schedule-must-avoid-minute-zero -- each run exports the interval that just closed, and batch_export.jitter already spreads the start
intervals=[ScheduleIntervalSpec(every=batch_export.interval_time_delta)],
jitter=batch_export.jitter,
time_zone_name=timezone,
Expand Down
1 change: 1 addition & 0 deletions products/today/backend/temporal/schedule.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ async def create_today_briefing_schedule(client: Client) -> None:
task_queue=settings.GENERAL_PURPOSE_TASK_QUEUE,
execution_timeout=timedelta(minutes=SCHEDULE_WINDOW_MINUTES),
),
# nosemgrep: schedule-must-avoid-minute-zero -- each run writes the briefings that start in the next window, so an offset shortens how far ahead of 8:00 local a briefing is written
spec=ScheduleSpec(intervals=[ScheduleIntervalSpec(every=timedelta(minutes=SCHEDULE_WINDOW_MINUTES))]),
policy=SchedulePolicy(overlap=ScheduleOverlapPolicy.SKIP),
)
Expand Down
Loading