Skip to content

Commit af61f95

Browse files
authored
BUG: Fix IDEX add_l1a_partition sensor (#1527)
* fix idex partition * add test * add comments to tests
1 parent 364b7c6 commit af61f95

2 files changed

Lines changed: 64 additions & 10 deletions

File tree

sds_data_manager/orchestration/custom_partitions.py

Lines changed: 8 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -138,8 +138,9 @@ def add_daily_partitions(context: SensorEvaluationContext):
138138
@sensor(minimum_interval_seconds=86400)
139139
def add_idex_10_day_partitions(context: SensorEvaluationContext):
140140
"""Alert Dagster when new IDEX 10-day partitions should be made."""
141-
start_date = context.cursor or config.MISSION_START_TIME
142-
start_dt = datetime.datetime.fromisoformat(start_date).replace(
141+
# These partitions come from a static cadence CSV, so we always evaluate from
142+
# mission start to avoid cursor drift shrinking the effective window.
143+
start_dt = datetime.datetime.fromisoformat(config.MISSION_START_TIME).replace(
143144
tzinfo=datetime.timezone.utc
144145
)
145146
end_dt = datetime.datetime.now(datetime.timezone.utc) + datetime.timedelta(
@@ -187,9 +188,7 @@ def add_idex_10_day_partitions(context: SensorEvaluationContext):
187188
partition_requests.append(idex10_partitions.build_add_request(partition_names))
188189
context.log.info(f"Registered new dynamic partitions: {partition_names}")
189190

190-
return SensorResult(
191-
dynamic_partitions_requests=partition_requests, cursor=end_dt.isoformat()
192-
)
191+
return SensorResult(dynamic_partitions_requests=partition_requests)
193192

194193

195194
##### THIS TELLS DAGSTER THAT SOME FILES ARE DIVIDED UP BY 30-day
@@ -199,8 +198,9 @@ def add_idex_10_day_partitions(context: SensorEvaluationContext):
199198
@sensor(minimum_interval_seconds=86400)
200199
def add_idex_30_day_partitions(context: SensorEvaluationContext):
201200
"""Alert Dagster when new IDEX 30-day partitions should be made."""
202-
start_date = context.cursor or config.MISSION_START_TIME
203-
start_dt = datetime.datetime.fromisoformat(start_date).replace(
201+
# These partitions come from a static cadence CSV, so we always evaluate from
202+
# mission start to avoid cursor drift shrinking the effective window.
203+
start_dt = datetime.datetime.fromisoformat(config.MISSION_START_TIME).replace(
204204
tzinfo=datetime.timezone.utc
205205
)
206206
end_dt = datetime.datetime.now(datetime.timezone.utc) + datetime.timedelta(
@@ -252,9 +252,7 @@ def add_idex_30_day_partitions(context: SensorEvaluationContext):
252252
partition_requests.append(idex30_partitions.build_add_request(partition_names))
253253
context.log.info(f"Registered new dynamic partitions: {partition_names}")
254254

255-
return SensorResult(
256-
dynamic_partitions_requests=partition_requests, cursor=end_dt.isoformat()
257-
)
255+
return SensorResult(dynamic_partitions_requests=partition_requests)
258256

259257

260258
# Run daily (24 hours = 86400 seconds)
Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,56 @@
1+
"""Tests for dynamic partition sensors in custom_partitions."""
2+
3+
import datetime
4+
from unittest.mock import patch
5+
6+
from dagster import build_sensor_context, instance_for_test
7+
8+
from sds_data_manager.orchestration import custom_partitions
9+
10+
11+
# wrap the real datetime *module*; override only datetime.datetime.now
12+
class _FrozenDatetime(datetime.datetime):
13+
@classmethod
14+
def now(cls, tz=None):
15+
return datetime.datetime(2025, 9, 29, tzinfo=datetime.timezone.utc)
16+
17+
18+
@patch.object(custom_partitions.datetime, "datetime", _FrozenDatetime)
19+
def test_add_idex_10_day_partitions():
20+
"""Check that add_idex_10_day_partitions adds the correct partitions."""
21+
with instance_for_test() as instance:
22+
# Mock existing partitions
23+
instance.add_dynamic_partitions(
24+
"idex_10_day_partitions",
25+
["idex10_2025-09-27T00:00:00_to_2025-10-07T00:00:00"],
26+
)
27+
context = build_sensor_context(instance=instance)
28+
# Trigger the sensor. This should add more partitions.
29+
sensor_result = custom_partitions.add_idex_10_day_partitions(context)
30+
31+
new_partitions = sensor_result.dynamic_partitions_requests[0].partition_keys
32+
assert new_partitions == [
33+
"idex10_2025-10-07T00:00:00_to_2025-10-17T00:00:00",
34+
"idex10_2025-10-17T00:00:00_to_2025-10-27T00:00:00",
35+
"idex10_2025-10-27T00:00:00_to_2025-11-06T00:00:00",
36+
]
37+
38+
39+
@patch.object(custom_partitions.datetime, "datetime", _FrozenDatetime)
40+
def test_add_idex_30_day_partitions():
41+
"""Check that add_idex_30_day_partitions adds the correct partitions."""
42+
with instance_for_test() as instance:
43+
# Mock existing partitions
44+
instance.add_dynamic_partitions(
45+
"idex_30_day_partitions",
46+
["idex30_2025-09-27T00:00:00_to_2025-10-07T00:00:00"],
47+
)
48+
context = build_sensor_context(instance=instance)
49+
# Trigger the sensor. This should add more partitions.
50+
sensor_result = custom_partitions.add_idex_30_day_partitions(context)
51+
52+
new_partitions = sensor_result.dynamic_partitions_requests[0].partition_keys
53+
assert new_partitions == [
54+
"idex30_2025-09-24T00:00:00_to_2025-09-27T00:00:00",
55+
"idex30_2025-09-27T00:00:00_to_2025-10-27T00:00:00",
56+
]

0 commit comments

Comments
 (0)