Skip to content

Commit 9713a23

Browse files
tmplummerCopilot
andauthored
1510 add reconstructed ephemeris to kernels that trigger only new coverage (#1512)
* Handle reconstructed SPICE coverage triggers * Clarify predicted SPICE trigger filtering * Explain growing SPICE kernels * Add predicts back to event bridge kernels * Update test with correct spice event list --------- Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com>
1 parent 675c32f commit 9713a23

5 files changed

Lines changed: 195 additions & 12 deletions

File tree

sds_data_manager/lambda_code/SDSCode/pipeline_lambdas/spice_indexer.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -519,6 +519,7 @@ def send_spice_event(spice_obj: SPICEFilePath, s3_key: str):
519519
spice_events = [
520520
"attitude_history",
521521
"attitude_predict",
522+
"pointing_attitude",
522523
"ephemeris_reconstructed",
523524
"ephemeris_nominal",
524525
"ephemeris_predict",

sds_data_manager/orchestration/imap_job.py

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -634,7 +634,11 @@ def trigger_from_new_non_science_inputs(
634634
if new_files:
635635
latest_ingestion_date = max(f.ingestion_date for f in new_files)
636636

637-
if cursor_str == config.MISSION_START_TIME:
637+
if dependency.source in spice.NON_TRIGGERING_KERNEL_TYPES:
638+
context.log.info(
639+
f"Skipping trigger evaluation for {dependency.source}."
640+
)
641+
elif cursor_str == config.MISSION_START_TIME:
638642
# Just process everything, don't bother looping through all new files.
639643
partitions = dagster_utilities.get_affected_partitions(
640644
context,

sds_data_manager/orchestration/spice.py

Lines changed: 27 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,25 @@
1414
logger.setLevel(logging.INFO)
1515

1616

17+
# These kernels are delivered as append-only time series: newer files either
18+
# extend an existing coverage window or replace it with a higher version of the
19+
# same window, so downstream triggering should be narrowed to genuinely new
20+
# coverage where possible.
21+
GROWING_KERNEL_TYPES = (
22+
"attitude_history",
23+
"pointing_attitude",
24+
"ephemeris_reconstructed",
25+
)
26+
# Kernel types that we don't want to trigger processing.
27+
# Predicted ephemeris appears as `ephemeris_predicted` in dependency YAML and
28+
# as `ephemeris_predict` in indexed SPICE metadata, so guard against both.
29+
NON_TRIGGERING_KERNEL_TYPES = (
30+
"attitude_predict",
31+
"ephemeris_predict",
32+
"ephemeris_predicted",
33+
)
34+
35+
1736
def check_requested_kernels(combined_kernel_sources, metakernel_files):
1837
"""Check if all requested kernels are present in the metakernel files.
1938
@@ -197,7 +216,14 @@ def _parse_interval_list(
197216
if not raw_intervals:
198217
return []
199218
return [
200-
[datetime.datetime.fromisoformat(start), datetime.datetime.fromisoformat(end)]
219+
[
220+
start
221+
if isinstance(start, datetime.datetime)
222+
else datetime.datetime.fromisoformat(start),
223+
end
224+
if isinstance(end, datetime.datetime)
225+
else datetime.datetime.fromisoformat(end),
226+
]
201227
for start, end in raw_intervals
202228
]
203229

@@ -269,15 +295,6 @@ def subtract_intervals(
269295
return leftover
270296

271297

272-
# Kernel types whose coverage grows by appending segments to the same file
273-
# series over time (attitude_history: MOC-delivered AH kernels; pointing_attitude:
274-
# the DPS kernel produced by the spacecraft l1a pointing-attitude job). Both use
275-
# the same start-date_end-date_version filename convention and the same
276-
# ingestion-time SPICE segment decomposition, so the same trigger-range rules
277-
# apply to both.
278-
GROWING_KERNEL_TYPES = ("attitude_history", "pointing_attitude")
279-
280-
281298
def get_growing_kernel_trigger_ranges(
282299
session, kernel_type: str, new_files: list["models.SPICEFiles"]
283300
) -> list[tuple[datetime.datetime, datetime.datetime]]:
Lines changed: 141 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,141 @@
1+
"""Tests for SPICE dependency trigger narrowing."""
2+
3+
from datetime import datetime
4+
from types import SimpleNamespace
5+
from unittest.mock import MagicMock, patch
6+
7+
from dagster import DagsterInstance, build_sensor_context
8+
9+
from sds_data_manager.lambda_code.SDSCode.database import models
10+
from sds_data_manager.orchestration.dependency import DependencyConfigReader
11+
from sds_data_manager.orchestration.imap_job import IMAPJobHandler
12+
13+
14+
def _mag_job_handler():
15+
"""Return a job handler with ephemeris dependencies."""
16+
reader = DependencyConfigReader()
17+
return IMAPJobHandler(reader.config[("mag", "l1d", "norm-srf")])
18+
19+
20+
def _session_returning(new_files, predecessor):
21+
"""Return a mock session for trigger range testing."""
22+
session = MagicMock()
23+
session.query.return_value.filter.return_value.all.return_value = new_files
24+
first = session.query.return_value.filter.return_value.order_by.return_value.first
25+
first.return_value = predecessor
26+
return session
27+
28+
29+
def _spice_file(file_name, min_dt, max_dt, intervals):
30+
"""Build a lightweight SPICE file object for tests."""
31+
return SimpleNamespace(
32+
file_name=file_name,
33+
min_date_datetime=min_dt,
34+
max_date_datetime=max_dt,
35+
file_intervals_datetime=intervals,
36+
ingestion_date=datetime(2025, 1, 5),
37+
)
38+
39+
40+
def test_ephemeris_reconstructed_only_triggers_new_daily_partitions():
41+
"""Ephemeris reconstructed kernels should only trigger for new coverage."""
42+
instance = DagsterInstance.ephemeral()
43+
instance.add_dynamic_partitions(
44+
"daily_partitions",
45+
[
46+
"daily_2025-01-01T00:00:00_to_2025-01-02T00:00:00",
47+
"daily_2025-01-02T00:00:00_to_2025-01-03T00:00:00",
48+
"daily_2025-01-03T00:00:00_to_2025-01-04T00:00:00",
49+
],
50+
)
51+
context = build_sensor_context(instance=instance)
52+
53+
job_handler = _mag_job_handler()
54+
dependency = next(
55+
dep
56+
for dep in job_handler.job_config.spice_inputs
57+
if dep.source == "ephemeris_reconstructed"
58+
)
59+
60+
min_dt = datetime(2025, 1, 1)
61+
old_max = datetime(2025, 1, 2)
62+
new_max = datetime(2025, 1, 4)
63+
predecessor = _spice_file(
64+
"imap_2025_001_2025_002_01.bsp",
65+
min_dt,
66+
old_max,
67+
[[min_dt, old_max]],
68+
)
69+
new_file = _spice_file(
70+
"imap_2025_001_2025_004_01.bsp",
71+
min_dt,
72+
new_max,
73+
[[min_dt, old_max], [old_max, new_max]],
74+
)
75+
session = _session_returning([new_file], predecessor)
76+
cursors = {dependency.to_dagster_name(): "2024-01-01T00:00:00"}
77+
78+
with patch(
79+
"sds_data_manager.orchestration.imap_job.db.Session"
80+
) as mock_session_cls:
81+
mock_session_cls.return_value.__enter__.return_value = session
82+
mock_session_cls.return_value.__exit__.return_value = False
83+
partitions = job_handler.trigger_from_new_non_science_inputs(
84+
context,
85+
dependency,
86+
cursors,
87+
models.SPICEFiles,
88+
models.SPICEFiles.kernel_type,
89+
None,
90+
"min_date_datetime",
91+
"max_date_datetime",
92+
)
93+
94+
assert set(partitions) == {
95+
"daily_2025-01-02T00:00:00_to_2025-01-03T00:00:00",
96+
"daily_2025-01-03T00:00:00_to_2025-01-04T00:00:00",
97+
}
98+
99+
100+
def test_predicted_ephemeris_does_not_trigger_partitions():
101+
"""Predicted ephemeris kernels should not trigger reprocessing."""
102+
instance = DagsterInstance.ephemeral()
103+
instance.add_dynamic_partitions(
104+
"daily_partitions",
105+
["daily_2025-01-01T00:00:00_to_2025-01-02T00:00:00"],
106+
)
107+
context = build_sensor_context(instance=instance)
108+
109+
job_handler = _mag_job_handler()
110+
dependency = next(
111+
dep
112+
for dep in job_handler.job_config.spice_inputs
113+
if dep.source == "ephemeris_predicted"
114+
)
115+
116+
new_file = _spice_file(
117+
"imap_2025_001_2025_002_01_pred.bsp",
118+
datetime(2025, 1, 1),
119+
datetime(2025, 1, 2),
120+
[[datetime(2025, 1, 1), datetime(2025, 1, 2)]],
121+
)
122+
session = _session_returning([new_file], None)
123+
cursors = {dependency.to_dagster_name(): "2024-01-01T00:00:00"}
124+
125+
with patch(
126+
"sds_data_manager.orchestration.imap_job.db.Session"
127+
) as mock_session_cls:
128+
mock_session_cls.return_value.__enter__.return_value = session
129+
mock_session_cls.return_value.__exit__.return_value = False
130+
partitions = job_handler.trigger_from_new_non_science_inputs(
131+
context,
132+
dependency,
133+
cursors,
134+
models.SPICEFiles,
135+
models.SPICEFiles.kernel_type,
136+
None,
137+
"min_date_datetime",
138+
"max_date_datetime",
139+
)
140+
141+
assert partitions == []

tests/lambda_endpoints/test_spice_indexer_lambda.py

Lines changed: 21 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@
44
import os
55
from datetime import datetime, timedelta
66
from pathlib import Path
7-
from unittest.mock import patch
7+
from unittest.mock import Mock, patch
88

99
import numpy as np
1010
import pytest
@@ -456,6 +456,26 @@ def test_send_spice_event(session, events_client, s3_client):
456456
spice_indexer.lambda_handler(event, None)
457457

458458

459+
def test_send_spice_event_filters_kernel_types(events_client):
460+
"""Only reconstructed/current coverage kernels should emit events."""
461+
with patch(
462+
"sds_data_manager.lambda_code.SDSCode.pipeline_lambdas.spice_indexer.boto3.client",
463+
return_value=events_client,
464+
):
465+
for kernel_type in (
466+
"attitude_history",
467+
"pointing_attitude",
468+
"ephemeris_reconstructed",
469+
"attitude_predict",
470+
"ephemeris_predict",
471+
):
472+
result = spice_indexer.send_spice_event(
473+
Mock(spice_metadata={"type": kernel_type}),
474+
"imap/spice/test/file",
475+
)
476+
assert result["ResponseMetadata"]["HTTPStatusCode"] == 200
477+
478+
459479
@patch(
460480
"sds_data_manager.lambda_code.SDSCode.pipeline_lambdas.spice_indexer.download_from_s3"
461481
)

0 commit comments

Comments
 (0)