From 5df7bbd10afaf2abeff25ecbda943c71226e53c8 Mon Sep 17 00:00:00 2001 From: Moira Andrews Date: Wed, 18 Feb 2026 13:51:48 -0800 Subject: [PATCH 01/24] add a check and update for observation status --- tom_observations/cadences/retry_failed_observations.py | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/tom_observations/cadences/retry_failed_observations.py b/tom_observations/cadences/retry_failed_observations.py index fc15217bb..157069494 100644 --- a/tom_observations/cadences/retry_failed_observations.py +++ b/tom_observations/cadences/retry_failed_observations.py @@ -2,7 +2,7 @@ from dateutil.parser import parse from tom_observations.cadence import BaseCadenceForm, CadenceStrategy -from tom_observations.models import ObservationRecord +from tom_observations.models import ObservationRecord, DynamicCadence from tom_observations.facility import get_service_class @@ -23,6 +23,11 @@ class RetryFailedObservationsStrategy(CadenceStrategy): form = RetryFailedObservationsForm def run(self): + last_obs = self.dynamic_cadence.observation_group.observation_records.order_by('-created').first() + facility = get_service_class(last_obs.facility)() + facility.update_observation_status(last_obs.observation_id) # Updates the DB record + last_obs.refresh_from_db() + failed_observations = [obsr for obsr in self.dynamic_cadence.observation_group.observation_records.all() if obsr.failed] From 0c2c80dd1753ee242a20b27b5727d4411f926c60 Mon Sep 17 00:00:00 2001 From: Moira Andrews Date: Wed, 18 Feb 2026 13:52:12 -0800 Subject: [PATCH 02/24] update dynamic cadence on obs complete --- tom_observations/cadences/retry_failed_observations.py | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/tom_observations/cadences/retry_failed_observations.py b/tom_observations/cadences/retry_failed_observations.py index 157069494..413ff1c33 100644 --- a/tom_observations/cadences/retry_failed_observations.py +++ b/tom_observations/cadences/retry_failed_observations.py @@ -28,6 +28,12 @@ def run(self): facility.update_observation_status(last_obs.observation_id) # Updates the DB record last_obs.refresh_from_db() + if last_obs.status == 'COMPLETED': + obs_group = last_obs.observationgroup_set.first() + dynamic_cadence = DynamicCadence.objects.get(observation_group=obs_group) + dynamic_cadence.active = False + dynamic_cadence.save() + failed_observations = [obsr for obsr in self.dynamic_cadence.observation_group.observation_records.all() if obsr.failed] From 25298c3f7dc08e595b9db26f3e4b6a5571890027 Mon Sep 17 00:00:00 2001 From: Moira Andrews Date: Wed, 18 Feb 2026 13:52:54 -0800 Subject: [PATCH 03/24] remove excess whitespace --- tom_observations/cadences/retry_failed_observations.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tom_observations/cadences/retry_failed_observations.py b/tom_observations/cadences/retry_failed_observations.py index 413ff1c33..e10be522e 100644 --- a/tom_observations/cadences/retry_failed_observations.py +++ b/tom_observations/cadences/retry_failed_observations.py @@ -27,7 +27,7 @@ def run(self): facility = get_service_class(last_obs.facility)() facility.update_observation_status(last_obs.observation_id) # Updates the DB record last_obs.refresh_from_db() - + if last_obs.status == 'COMPLETED': obs_group = last_obs.observationgroup_set.first() dynamic_cadence = DynamicCadence.objects.get(observation_group=obs_group) From 982fe86b19f3a68e3fc5180a54e62ad906df1205 Mon Sep 17 00:00:00 2001 From: Joey Chatelain Date: Wed, 4 Mar 2026 15:29:53 -0700 Subject: [PATCH 04/24] mock status check --- tom_observations/tests/test_cadence.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/tom_observations/tests/test_cadence.py b/tom_observations/tests/test_cadence.py index 66cf49511..a106168ae 100644 --- a/tom_observations/tests/test_cadence.py +++ b/tom_observations/tests/test_cadence.py @@ -59,7 +59,9 @@ def setUp(self): cadence_strategy='Test Strategy', cadence_parameters={'cadence_frequency': 72}, active=True, observation_group=self.group) - def test_retry_when_failed_cadence(self, patch1, patch2, patch3, patch4): + @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', return_value={'state': 'CANCELED', + 'scheduled_start': None, 'scheduled_end': None}) + def test_retry_when_failed_cadence(self, patch1, patch2, patch3, patch4, mock_get_obs_status): num_records = self.group.observation_records.count() observing_record = self.group.observation_records.first() observing_record.status = 'CANCELED' From aed35bf7a9dd9cad0edd9aea4184aa320ca0a328 Mon Sep 17 00:00:00 2001 From: moira-andrews <85570657+moira-andrews@users.noreply.github.com> Date: Wed, 4 Mar 2026 15:07:02 -0800 Subject: [PATCH 05/24] Cleaning up logic and pulling changes from snex2 --- .../cadences/retry_failed_observations.py | 79 +++++++++++-------- 1 file changed, 46 insertions(+), 33 deletions(-) diff --git a/tom_observations/cadences/retry_failed_observations.py b/tom_observations/cadences/retry_failed_observations.py index e10be522e..0f6f25ec1 100644 --- a/tom_observations/cadences/retry_failed_observations.py +++ b/tom_observations/cadences/retry_failed_observations.py @@ -23,44 +23,57 @@ class RetryFailedObservationsStrategy(CadenceStrategy): form = RetryFailedObservationsForm def run(self): - last_obs = self.dynamic_cadence.observation_group.observation_records.order_by('-created').first() - facility = get_service_class(last_obs.facility)() - facility.update_observation_status(last_obs.observation_id) # Updates the DB record + records = self.dynamic_cadence.observation_group.observation_records.all().order_by('-created') + last_obs = records.first() + + if not last_obs: + return + + facility_class = get_service_class(last_obs.facility) + facility = facility_class() + facility.update_observation_status(last_obs.observation_id) last_obs.refresh_from_db() - if last_obs.status == 'COMPLETED': - obs_group = last_obs.observationgroup_set.first() - dynamic_cadence = DynamicCadence.objects.get(observation_group=obs_group) - dynamic_cadence.active = False - dynamic_cadence.save() + if not last_obs.terminal: + return + elif last_obs.status == 'COMPLETED': + self.dynamic_cadence.active = False + self.dynamic_cadence.save() + return + + if not last_obs.failed: + return + + observation_payload = last_obs.parameters.copy() - failed_observations = [obsr for obsr - in self.dynamic_cadence.observation_group.observation_records.all() - if obsr.failed] + start_keyword, end_keyword = facility.get_start_end_keywords() + observation_payload = self.advance_window( + observation_payload, start_keyword=start_keyword, end_keyword=end_keyword + ) + + obs_type = observation_payload.get('observation_type') + form = facility.get_form(obs_type)(observation_payload) + + if not form.is_valid(): + return + + observation_ids = facility.submit_observation(form.observation_payload()) new_observations = [] - for obs in failed_observations: - observation_payload = obs.parameters - facility = get_service_class(obs.facility)() - start_keyword, end_keyword = facility.get_start_end_keywords() - observation_payload = self.advance_window( - observation_payload, start_keyword=start_keyword, end_keyword=end_keyword + + for observation_id in observation_ids: + record = ObservationRecord.objects.create( + target=last_obs.target, + facility=facility.name, + parameters=observation_payload, + observation_id=observation_id ) - obs_type = obs.parameters.get('observation_type', None) - form = facility.get_form(obs_type)(data=observation_payload) - form.is_valid() - observation_ids = facility.submit_observation(form.observation_payload()) - - for observation_id in observation_ids: - # Create Observation record - record = ObservationRecord.objects.create( - target=obs.target, - facility=facility.name, - parameters=observation_payload, - observation_id=observation_id - ) - self.dynamic_cadence.observation_group.observation_records.add(record) - self.dynamic_cadence.observation_group.save() - new_observations.append(record) + self.dynamic_cadence.observation_group.observation_records.add(record) + new_observations.append(record) + + self.dynamic_cadence.observation_group.save() + + for obsr in new_observations: + facility.update_observation_status(obsr.observation_id) return new_observations From 834339d60dd62710d628df3c85b7881fa11632f2 Mon Sep 17 00:00:00 2001 From: Joey Chatelain Date: Wed, 4 Mar 2026 16:20:45 -0700 Subject: [PATCH 06/24] fix linting --- tom_observations/cadences/retry_failed_observations.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/tom_observations/cadences/retry_failed_observations.py b/tom_observations/cadences/retry_failed_observations.py index 0f6f25ec1..ecbf6ef82 100644 --- a/tom_observations/cadences/retry_failed_observations.py +++ b/tom_observations/cadences/retry_failed_observations.py @@ -2,7 +2,7 @@ from dateutil.parser import parse from tom_observations.cadence import BaseCadenceForm, CadenceStrategy -from tom_observations.models import ObservationRecord, DynamicCadence +from tom_observations.models import ObservationRecord from tom_observations.facility import get_service_class @@ -50,7 +50,7 @@ def run(self): observation_payload = self.advance_window( observation_payload, start_keyword=start_keyword, end_keyword=end_keyword ) - + obs_type = observation_payload.get('observation_type') form = facility.get_form(obs_type)(observation_payload) @@ -59,7 +59,7 @@ def run(self): observation_ids = facility.submit_observation(form.observation_payload()) new_observations = [] - + for observation_id in observation_ids: record = ObservationRecord.objects.create( target=last_obs.target, From c55dece2a16e9e106351aeeb7fd33cf164b1443c Mon Sep 17 00:00:00 2001 From: Joey Chatelain Date: Wed, 4 Mar 2026 16:22:44 -0700 Subject: [PATCH 07/24] missed one --- tom_observations/cadences/retry_failed_observations.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tom_observations/cadences/retry_failed_observations.py b/tom_observations/cadences/retry_failed_observations.py index ecbf6ef82..264ad736b 100644 --- a/tom_observations/cadences/retry_failed_observations.py +++ b/tom_observations/cadences/retry_failed_observations.py @@ -53,7 +53,7 @@ def run(self): obs_type = observation_payload.get('observation_type') form = facility.get_form(obs_type)(observation_payload) - + if not form.is_valid(): return From ee97c71c8ef9070be5f73930e5d91c753d4f5c16 Mon Sep 17 00:00:00 2001 From: Moira Andrews Date: Wed, 1 Apr 2026 12:09:07 -0400 Subject: [PATCH 08/24] update advance window to take the scheduled start and end times instead of the full cadence window start and end --- .../cadences/resume_cadence_after_failure.py | 11 ++++++++++- 1 file changed, 10 insertions(+), 1 deletion(-) diff --git a/tom_observations/cadences/resume_cadence_after_failure.py b/tom_observations/cadences/resume_cadence_after_failure.py index cfc663d8a..963a74bf8 100644 --- a/tom_observations/cadences/resume_cadence_after_failure.py +++ b/tom_observations/cadences/resume_cadence_after_failure.py @@ -47,9 +47,17 @@ def run(self): last_obs.refresh_from_db() # Gets the record updates # Boilerplate to get necessary properties for future calls - start_keyword, end_keyword = facility.get_start_end_keywords() observation_payload = last_obs.parameters + scheduled_start = last_obs.scheduled_start + scheduled_end = last_obs.scheduled_end + + # Add the scheduled start and end + observation_payload['scheduled_start'] = scheduled_start + observation_payload['scheduled_end'] = scheduled_end + start_keyword = 'scheduled_start' + end_keyword = 'scheduled_end' + # Cadence logic # If the observation hasn't finished, do nothing if not last_obs.terminal: @@ -95,6 +103,7 @@ def run(self): for obsr in new_observations: facility = get_service_class(obsr.facility)() facility.update_observation_status(obsr.observation_id) + obsr.refresh_from_db() # commit the updated observation status return new_observations From 56841536628fe952c9b83136317dd2fd68ef637d Mon Sep 17 00:00:00 2001 From: Moira Andrews Date: Wed, 1 Apr 2026 12:10:51 -0400 Subject: [PATCH 09/24] simplifying variables --- tom_observations/cadences/resume_cadence_after_failure.py | 7 ++----- 1 file changed, 2 insertions(+), 5 deletions(-) diff --git a/tom_observations/cadences/resume_cadence_after_failure.py b/tom_observations/cadences/resume_cadence_after_failure.py index 963a74bf8..6d01bd88e 100644 --- a/tom_observations/cadences/resume_cadence_after_failure.py +++ b/tom_observations/cadences/resume_cadence_after_failure.py @@ -49,12 +49,9 @@ def run(self): # Boilerplate to get necessary properties for future calls observation_payload = last_obs.parameters - scheduled_start = last_obs.scheduled_start - scheduled_end = last_obs.scheduled_end - # Add the scheduled start and end - observation_payload['scheduled_start'] = scheduled_start - observation_payload['scheduled_end'] = scheduled_end + observation_payload['scheduled_start'] = last_obs.scheduled_start + observation_payload['scheduled_end'] = last_obs.scheduled_end start_keyword = 'scheduled_start' end_keyword = 'scheduled_end' From 2382974ee80040a35ac05fe47487a4ab2be02bef Mon Sep 17 00:00:00 2001 From: Moira Andrews Date: Wed, 1 Apr 2026 13:43:26 -0400 Subject: [PATCH 10/24] Updated logic to be a shorter window --- .../cadences/resume_cadence_after_failure.py | 20 ++++++++++++------- 1 file changed, 13 insertions(+), 7 deletions(-) diff --git a/tom_observations/cadences/resume_cadence_after_failure.py b/tom_observations/cadences/resume_cadence_after_failure.py index 6d01bd88e..44868f4b6 100644 --- a/tom_observations/cadences/resume_cadence_after_failure.py +++ b/tom_observations/cadences/resume_cadence_after_failure.py @@ -43,17 +43,15 @@ def run(self): # Make a call to the facility to get the current status of the observation facility = get_service_class(last_obs.facility)() + start_keyword, end_keyword = facility.get_start_end_keywords() facility.update_observation_status(last_obs.observation_id) # Updates the DB record last_obs.refresh_from_db() # Gets the record updates # Boilerplate to get necessary properties for future calls observation_payload = last_obs.parameters - # Add the scheduled start and end - observation_payload['scheduled_start'] = last_obs.scheduled_start + # Only needs the scheduled end from the observatory observation_payload['scheduled_end'] = last_obs.scheduled_end - start_keyword = 'scheduled_start' - end_keyword = 'scheduled_end' # Cadence logic # If the observation hasn't finished, do nothing @@ -61,7 +59,11 @@ def run(self): return elif last_obs.failed: # If the observation failed # Submit next observation to be taken as soon as possible with the same window length - window_length = parse(observation_payload[end_keyword]) - parse(observation_payload[start_keyword]) + # Make the window length the cadence frequency or 24 hours, whichever is shorter. + cadence_frequency = self.dynamic_cadence.cadence_parameters.get('cadence_frequency') + if not cadence_frequency: + raise Exception(f'The {self.name} strategy requires a cadence_frequency cadence_parameter.') + window_length = 24 if cadence_frequency > 24 else cadence_frequency observation_payload[start_keyword] = datetime.now().isoformat() observation_payload[end_keyword] = (parse(observation_payload[start_keyword]) + window_length).isoformat() else: # If the observation succeeded @@ -109,9 +111,13 @@ def advance_window(self, observation_payload, start_keyword='start', end_keyword if not cadence_frequency: raise Exception(f'The {self.name} strategy requires a cadence_frequency cadence_parameter.') advance_window_hours = cadence_frequency - window_length = parse(observation_payload[end_keyword]) - parse(observation_payload[start_keyword]) - new_start = parse(observation_payload[start_keyword]) + timedelta(hours=advance_window_hours) + # Window length for the observation should be every 24 hours unless the frequency is less than 24 + # then just the cadence frequency + window_length = 24 if cadence_frequency > 24 else cadence_frequency + + # Define new start to be at the end of the previous observation + cadence in hours + new_start = parse(observation_payload['scheduled_end']) + timedelta(hours=advance_window_hours) if new_start < datetime.now(): # Ensure that the new window isn't in the past new_start = datetime.now() new_end = new_start + window_length From 4df78e9bb473fdb75ed114e74b857f80927bb561 Mon Sep 17 00:00:00 2001 From: Moira Andrews Date: Wed, 1 Apr 2026 17:03:36 -0400 Subject: [PATCH 11/24] changed 24 hour to a settings variable call --- .../cadences/resume_cadence_after_failure.py | 9 +++++++-- .../cadences/retry_failed_observations.py | 11 ++++++++--- 2 files changed, 15 insertions(+), 5 deletions(-) diff --git a/tom_observations/cadences/resume_cadence_after_failure.py b/tom_observations/cadences/resume_cadence_after_failure.py index 44868f4b6..f03b6eaa5 100644 --- a/tom_observations/cadences/resume_cadence_after_failure.py +++ b/tom_observations/cadences/resume_cadence_after_failure.py @@ -1,5 +1,6 @@ from datetime import datetime, timedelta from dateutil.parser import parse +from django.conf import settings import logging from tom_observations.cadence import BaseCadenceForm, CadenceStrategy @@ -113,8 +114,12 @@ def advance_window(self, observation_payload, start_keyword='start', end_keyword advance_window_hours = cadence_frequency # Window length for the observation should be every 24 hours unless the frequency is less than 24 - # then just the cadence frequency - window_length = 24 if cadence_frequency > 24 else cadence_frequency + # then just the cadence frequency (24 hour window defined in the settings) + if settings.OBS_WINDOW_MINIMUM: + min_window = settings.OBS_WINDOW_MINIMUM + else: + min_window = 24 + window_length = min_window if cadence_frequency > min_window else cadence_frequency # Define new start to be at the end of the previous observation + cadence in hours new_start = parse(observation_payload['scheduled_end']) + timedelta(hours=advance_window_hours) diff --git a/tom_observations/cadences/retry_failed_observations.py b/tom_observations/cadences/retry_failed_observations.py index fc15217bb..fc4845004 100644 --- a/tom_observations/cadences/retry_failed_observations.py +++ b/tom_observations/cadences/retry_failed_observations.py @@ -1,5 +1,6 @@ from datetime import timedelta from dateutil.parser import parse +from django.conf import settings from tom_observations.cadence import BaseCadenceForm, CadenceStrategy from tom_observations.models import ObservationRecord @@ -57,9 +58,13 @@ def advance_window(self, observation_payload, start_keyword='start', end_keyword cadence_frequency = self.dynamic_cadence.cadence_parameters.get('cadence_frequency') if not cadence_frequency: raise Exception(f'The {self.name} strategy requires a cadence_frequency cadence_parameter.') - advance_window_hours = cadence_frequency - new_start = parse(observation_payload[start_keyword]) + timedelta(hours=advance_window_hours) - new_end = parse(observation_payload[end_keyword]) + timedelta(hours=advance_window_hours) + if settings.OBS_WINDOW_MINIMUM: + min_window = settings.OBS_WINDOW_MINIMUM + else: + min_window = 24 + window_length = min_window if cadence_frequency > min_window else cadence_frequency + new_start = parse(observation_payload[start_keyword]) + timedelta(hours=window_length) + new_end = parse(observation_payload[end_keyword]) + timedelta(hours=window_length) observation_payload[start_keyword] = new_start.isoformat() observation_payload[end_keyword] = new_end.isoformat() From 666c48d6185e3c547e94ad9147d4b9d082d93a89 Mon Sep 17 00:00:00 2001 From: Moira Andrews Date: Wed, 1 Apr 2026 17:23:10 -0400 Subject: [PATCH 12/24] updated logic and handling to be more direct and have runcadencestrategies display expected response --- .../cadences/retry_failed_observations.py | 78 ++++++++++--------- .../commands/runcadencestrategies.py | 4 + 2 files changed, 45 insertions(+), 37 deletions(-) diff --git a/tom_observations/cadences/retry_failed_observations.py b/tom_observations/cadences/retry_failed_observations.py index 264ad736b..a65640bf4 100644 --- a/tom_observations/cadences/retry_failed_observations.py +++ b/tom_observations/cadences/retry_failed_observations.py @@ -1,10 +1,13 @@ from datetime import timedelta from dateutil.parser import parse +import logging from tom_observations.cadence import BaseCadenceForm, CadenceStrategy -from tom_observations.models import ObservationRecord +from tom_observations.models import ObservationRecord, DynamicCadence from tom_observations.facility import get_service_class +logger = logging.getLogger(__name__) + class RetryFailedObservationsForm(BaseCadenceForm): pass @@ -31,51 +34,52 @@ def run(self): facility_class = get_service_class(last_obs.facility) facility = facility_class() + start_keyword, end_keyword = facility.get_start_end_keywords() facility.update_observation_status(last_obs.observation_id) last_obs.refresh_from_db() - if not last_obs.terminal: + if not last_obs.terminal: #observation is still pending, do nothing return - elif last_obs.status == 'COMPLETED': + + elif not last_obs.failed: #observation succeeded self.dynamic_cadence.active = False self.dynamic_cadence.save() - return - - if not last_obs.failed: - return - - observation_payload = last_obs.parameters.copy() - - start_keyword, end_keyword = facility.get_start_end_keywords() - observation_payload = self.advance_window( - observation_payload, start_keyword=start_keyword, end_keyword=end_keyword - ) - - obs_type = observation_payload.get('observation_type') - form = facility.get_form(obs_type)(observation_payload) - - if not form.is_valid(): - return + return 'COMPLETED' - observation_ids = facility.submit_observation(form.observation_payload()) - new_observations = [] + else: #observation failed, submit a new one + observation_payload = last_obs.parameters.copy() - for observation_id in observation_ids: - record = ObservationRecord.objects.create( - target=last_obs.target, - facility=facility.name, - parameters=observation_payload, - observation_id=observation_id + observation_payload = self.advance_window( + observation_payload, start_keyword=start_keyword, end_keyword=end_keyword ) - self.dynamic_cadence.observation_group.observation_records.add(record) - new_observations.append(record) - - self.dynamic_cadence.observation_group.save() - - for obsr in new_observations: - facility.update_observation_status(obsr.observation_id) - - return new_observations + + obs_type = observation_payload.get('observation_type') + form = facility.get_form(obs_type)(observation_payload) + + if not form.is_valid(): + logger.error(msg=f'Unable to submit next cadenced observation: {form.errors}') + raise Exception(f'Unable to submit next cadenced observation: {form.errors}') + + observation_ids = facility.submit_observation(form.observation_payload()) + new_observations = [] + + for observation_id in observation_ids: + record = ObservationRecord.objects.create( + target=last_obs.target, + facility=facility.name, + parameters=observation_payload, + observation_id=observation_id + ) + self.dynamic_cadence.observation_group.observation_records.add(record) + new_observations.append(record) + + self.dynamic_cadence.observation_group.save() + + for obsr in new_observations: + facility.update_observation_status(obsr.observation_id) + obsr.refresh_from_db() + + return new_observations def advance_window(self, observation_payload, start_keyword='start', end_keyword='end'): cadence_frequency = self.dynamic_cadence.cadence_parameters.get('cadence_frequency') diff --git a/tom_observations/management/commands/runcadencestrategies.py b/tom_observations/management/commands/runcadencestrategies.py index 73adbbec8..2ce50529e 100644 --- a/tom_observations/management/commands/runcadencestrategies.py +++ b/tom_observations/management/commands/runcadencestrategies.py @@ -36,6 +36,10 @@ def handle(self, *args, **kwargs): continue if not new_observations: logger.log(msg=f'No changes from dynamic cadence {cg}', level=logging.INFO) + elif new_observations == 'COMPLETED': + logger.log(msg=f'''Single observation obtained for {cg}, + no new observation submitted.''', + level=logging.INFO) else: logger.log(msg=f'''Cadence update completed for dynamic cadence {cg}, {len(new_observations)} new observations created.''', From 52743dc1ac8bf756657d9796e396a93619a6238f Mon Sep 17 00:00:00 2001 From: Moira Andrews Date: Wed, 1 Apr 2026 17:33:24 -0400 Subject: [PATCH 13/24] updated testing to follow the resume cadence strategy of setting no validation errors, added a test for a completed observation to turn off the cadence --- tom_observations/tests/test_cadence.py | 20 +++++++++++++++++++- 1 file changed, 19 insertions(+), 1 deletion(-) diff --git a/tom_observations/tests/test_cadence.py b/tom_observations/tests/test_cadence.py index a106168ae..e6e6ee207 100644 --- a/tom_observations/tests/test_cadence.py +++ b/tom_observations/tests/test_cadence.py @@ -61,7 +61,8 @@ def setUp(self): @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', return_value={'state': 'CANCELED', 'scheduled_start': None, 'scheduled_end': None}) - def test_retry_when_failed_cadence(self, patch1, patch2, patch3, patch4, mock_get_obs_status): + def test_retry_when_failed_cadence_failed_obs(self, patch1, patch2, patch3, patch4, mock_get_obs_status, mock_validate_obs): + mock_validate_obs.return_value = {} num_records = self.group.observation_records.count() observing_record = self.group.observation_records.first() observing_record.status = 'CANCELED' @@ -78,6 +79,23 @@ def test_retry_when_failed_cadence(self, patch1, patch2, patch3, patch4, mock_ge parse(observing_record.parameters['start']), parse(new_records[0].parameters['start']) - timedelta(days=3) ) + + @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', return_value={'state': 'CANCELED', + 'scheduled_start': None, 'scheduled_end': None}) + def test_retry_when_failed_cadence_successful_obs(self, patch1, patch2, patch3, patch4, mock_get_obs_status, mock_validate_obs): + mock_validate_obs.return_value = {} + observing_record = self.group.observation_records.first() + observing_record.status = 'COMPLETE' + observing_record.save() + + strategy = RetryFailedObservationsStrategy(self.dynamic_cadence) + new_records = strategy.run() + self.group.refresh_from_db() + # Make sure the candence returned 'COMPLETED' + self.assertEqual(new_records, 'COMPLETED') + # Make sure the dynamic cadence was turned off + self.assertEqual(self.dynamic_cadence.active, False) + @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', return_value={'state': 'CANCELED', 'scheduled_start': None, 'scheduled_end': None}) From 6cab1e5c34e73b0882dac475a48de1fbc441b42e Mon Sep 17 00:00:00 2001 From: Moira Andrews Date: Wed, 1 Apr 2026 18:06:20 -0400 Subject: [PATCH 14/24] adds the window minimum to the form clean function --- tom_observations/facilities/lco.py | 17 +++++++++++++++-- 1 file changed, 15 insertions(+), 2 deletions(-) diff --git a/tom_observations/facilities/lco.py b/tom_observations/facilities/lco.py index 9e96b4972..784a7e686 100644 --- a/tom_observations/facilities/lco.py +++ b/tom_observations/facilities/lco.py @@ -883,7 +883,13 @@ def clean(self): """ cleaned_data = super().clean() start = cleaned_data.get('start') - cleaned_data['end'] = datetime.strftime(parse(start) + timedelta(hours=cleaned_data['cadence_frequency']), + if settings.OBS_WINDOW_MINIMUM: + min_window = settings.OBS_WINDOW_MINIMUM + else: + min_window = 24 + window_length = min_window if cleaned_data['cadence_frequency'] > min_window else cleaned_data['cadence_frequency'] + + cleaned_data['end'] = datetime.strftime(parse(start) + timedelta(hours=window_length), '%Y-%m-%dT%H:%M:%S') return cleaned_data @@ -1058,8 +1064,15 @@ def clean(self): """ cleaned_data = super().clean() cleaned_data['instrument_type'] = '2M0-FLOYDS-SCICAM' # SNEx only submits spectra to FLOYDS + start = cleaned_data.get('start') - cleaned_data['end'] = datetime.strftime(parse(start) + timedelta(hours=cleaned_data['cadence_frequency']), + if settings.OBS_WINDOW_MINIMUM: + min_window = settings.OBS_WINDOW_MINIMUM + else: + min_window = 24 + window_length = min_window if cleaned_data['cadence_frequency'] > min_window else cleaned_data['cadence_frequency'] + + cleaned_data['end'] = datetime.strftime(parse(start) + timedelta(hours=window_length), '%Y-%m-%dT%H:%M:%S') return cleaned_data From b9e8e366b63e4b034d2dddb2573442c3307f6b7c Mon Sep 17 00:00:00 2001 From: Moira Andrews Date: Thu, 2 Apr 2026 14:13:55 -0400 Subject: [PATCH 15/24] updated min window to be consistent and fall back to the cadence frequency --- .../cadences/resume_cadence_after_failure.py | 14 ++++++++++---- .../cadences/retry_failed_observations.py | 7 ++++--- 2 files changed, 14 insertions(+), 7 deletions(-) diff --git a/tom_observations/cadences/resume_cadence_after_failure.py b/tom_observations/cadences/resume_cadence_after_failure.py index f03b6eaa5..43509776c 100644 --- a/tom_observations/cadences/resume_cadence_after_failure.py +++ b/tom_observations/cadences/resume_cadence_after_failure.py @@ -64,7 +64,12 @@ def run(self): cadence_frequency = self.dynamic_cadence.cadence_parameters.get('cadence_frequency') if not cadence_frequency: raise Exception(f'The {self.name} strategy requires a cadence_frequency cadence_parameter.') - window_length = 24 if cadence_frequency > 24 else cadence_frequency + if settings.OBS_WINDOW_MINIMUM: + window_length = settings.OBS_WINDOW_MINIMUM + if window_length > cadence_frequency: + window_length = cadence_frequency + else: + window_length = cadence_frequency observation_payload[start_keyword] = datetime.now().isoformat() observation_payload[end_keyword] = (parse(observation_payload[start_keyword]) + window_length).isoformat() else: # If the observation succeeded @@ -116,10 +121,11 @@ def advance_window(self, observation_payload, start_keyword='start', end_keyword # Window length for the observation should be every 24 hours unless the frequency is less than 24 # then just the cadence frequency (24 hour window defined in the settings) if settings.OBS_WINDOW_MINIMUM: - min_window = settings.OBS_WINDOW_MINIMUM + window_length = settings.OBS_WINDOW_MINIMUM + if window_length > cadence_frequency: + window_length = cadence_frequency else: - min_window = 24 - window_length = min_window if cadence_frequency > min_window else cadence_frequency + window_length = cadence_frequency # Define new start to be at the end of the previous observation + cadence in hours new_start = parse(observation_payload['scheduled_end']) + timedelta(hours=advance_window_hours) diff --git a/tom_observations/cadences/retry_failed_observations.py b/tom_observations/cadences/retry_failed_observations.py index fc4845004..76833306d 100644 --- a/tom_observations/cadences/retry_failed_observations.py +++ b/tom_observations/cadences/retry_failed_observations.py @@ -59,10 +59,11 @@ def advance_window(self, observation_payload, start_keyword='start', end_keyword if not cadence_frequency: raise Exception(f'The {self.name} strategy requires a cadence_frequency cadence_parameter.') if settings.OBS_WINDOW_MINIMUM: - min_window = settings.OBS_WINDOW_MINIMUM + window_length = settings.OBS_WINDOW_MINIMUM + if window_length > cadence_frequency: + window_length = cadence_frequency else: - min_window = 24 - window_length = min_window if cadence_frequency > min_window else cadence_frequency + window_length = cadence_frequency new_start = parse(observation_payload[start_keyword]) + timedelta(hours=window_length) new_end = parse(observation_payload[end_keyword]) + timedelta(hours=window_length) observation_payload[start_keyword] = new_start.isoformat() From 2d858ceac66d6072e68b69dc675273744b820658 Mon Sep 17 00:00:00 2001 From: Moira Andrews Date: Thu, 2 Apr 2026 14:16:20 -0400 Subject: [PATCH 16/24] updated window length call in facility --- tom_observations/facilities/lco.py | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/tom_observations/facilities/lco.py b/tom_observations/facilities/lco.py index 784a7e686..4f4a6897d 100644 --- a/tom_observations/facilities/lco.py +++ b/tom_observations/facilities/lco.py @@ -883,11 +883,13 @@ def clean(self): """ cleaned_data = super().clean() start = cleaned_data.get('start') + cadence_frequency = cleaned_data['cadence_frequency'] if settings.OBS_WINDOW_MINIMUM: - min_window = settings.OBS_WINDOW_MINIMUM + window_length = settings.OBS_WINDOW_MINIMUM + if window_length > cadence_frequency: + window_length = cadence_frequency else: - min_window = 24 - window_length = min_window if cleaned_data['cadence_frequency'] > min_window else cleaned_data['cadence_frequency'] + window_length = cadence_frequency cleaned_data['end'] = datetime.strftime(parse(start) + timedelta(hours=window_length), '%Y-%m-%dT%H:%M:%S') From e081cf018f74af0e456d57a8edb99c185cc2717d Mon Sep 17 00:00:00 2001 From: Moira Andrews Date: Wed, 24 Jun 2026 22:36:34 -0700 Subject: [PATCH 17/24] updated functions to match what is used in SNEx for sequencing --- tom_base/settings.py | 3 + .../cadences/resume_cadence_after_failure.py | 88 +++++--- .../cadences/retry_failed_observations.py | 208 +++++++++++++----- tom_observations/facilities/lco.py | 39 ++-- 4 files changed, 228 insertions(+), 110 deletions(-) diff --git a/tom_base/settings.py b/tom_base/settings.py index 35bc64fe8..88d57f033 100644 --- a/tom_base/settings.py +++ b/tom_base/settings.py @@ -285,9 +285,12 @@ TOM_CADENCE_STRATEGIES = [ 'tom_observations.cadences.retry_failed_observations.RetryFailedObservationsStrategy', + 'tom_observations.cadences.retry_failed_observations.RetryUntilDeadlineStrategy', 'tom_observations.cadences.resume_cadence_after_failure.ResumeCadenceAfterFailureStrategy' ] +OBS_WINDOW_MINIMUM = 24 + # Define extra target fields here. Types can be any of "number", "string", "boolean" or "datetime" # See https://tomtoolkit.github.io/docs/target_fields for documentation on this feature # For example: diff --git a/tom_observations/cadences/resume_cadence_after_failure.py b/tom_observations/cadences/resume_cadence_after_failure.py index 43509776c..e480ada11 100644 --- a/tom_observations/cadences/resume_cadence_after_failure.py +++ b/tom_observations/cadences/resume_cadence_after_failure.py @@ -2,6 +2,7 @@ from dateutil.parser import parse from django.conf import settings import logging +from django.utils import timezone from tom_observations.cadence import BaseCadenceForm, CadenceStrategy from tom_observations.models import ObservationRecord @@ -41,6 +42,8 @@ def update_observation_payload(self, observation_payload): def run(self): # gets the most recent observation because the next observation is just going to modify these parameters last_obs = self.dynamic_cadence.observation_group.observation_records.order_by('-created').first() + if not last_obs: + return # Make a call to the facility to get the current status of the observation facility = get_service_class(last_obs.facility)() @@ -48,30 +51,42 @@ def run(self): facility.update_observation_status(last_obs.observation_id) # Updates the DB record last_obs.refresh_from_db() # Gets the record updates - # Boilerplate to get necessary properties for future calls - observation_payload = last_obs.parameters - - # Only needs the scheduled end from the observatory - observation_payload['scheduled_end'] = last_obs.scheduled_end - # Cadence logic # If the observation hasn't finished, do nothing if not last_obs.terminal: return - elif last_obs.failed: # If the observation failed + + if last_obs.status == 'CANCELED': + logger.info(f'Observation {last_obs} was canceled, not resuming cadence') + return + + # Boilerplate to get necessary properties for future calls + observation_payload = last_obs.parameters.copy() + scheduled_end = last_obs.scheduled_end + if not scheduled_end: + logger.info(f'No observation end scheduled yet, falling back to end: {observation_payload[end_keyword]}') + scheduled_end = parse(observation_payload[end_keyword]) + + if isinstance(scheduled_end, str): + scheduled_end = parse(scheduled_end) + + if timezone.is_naive(scheduled_end): + scheduled_end = timezone.make_aware(scheduled_end) + + observation_payload['scheduled_end'] = scheduled_end.isoformat() + logger.info(f'Scheduled observation end: {scheduled_end}') + + if last_obs.failed: # If the observation failed # Submit next observation to be taken as soon as possible with the same window length - # Make the window length the cadence frequency or 24 hours, whichever is shorter. cadence_frequency = self.dynamic_cadence.cadence_parameters.get('cadence_frequency') - if not cadence_frequency: + if cadence_frequency is None: raise Exception(f'The {self.name} strategy requires a cadence_frequency cadence_parameter.') - if settings.OBS_WINDOW_MINIMUM: - window_length = settings.OBS_WINDOW_MINIMUM - if window_length > cadence_frequency: - window_length = cadence_frequency - else: - window_length = cadence_frequency - observation_payload[start_keyword] = datetime.now().isoformat() - observation_payload[end_keyword] = (parse(observation_payload[start_keyword]) + window_length).isoformat() + window_min = getattr(settings, 'OBS_WINDOW_MINIMUM', 24) + window_length = min(cadence_frequency, window_min) + now = timezone.now() + observation_payload[start_keyword] = now.isoformat() + observation_payload[end_keyword] = (now + timedelta(hours=window_length)).isoformat() + else: # If the observation succeeded # Advance window normally according to cadence parameters observation_payload = self.advance_window( @@ -83,10 +98,14 @@ def run(self): # Submission of the new observation to the facility obs_type = last_obs.parameters.get('observation_type') form = facility.get_form(obs_type)(data=observation_payload) + logger.info(f'obs payload: {observation_payload}') if form.is_valid(): observation_ids = facility.submit_observation(form.observation_payload()) else: - logger.error(msg=f'Unable to submit next cadenced observation: {form.errors}') + logger.error( + msg=f'Unable to submit next cadenced observation: {form.errors} ' + f'for ObservationRecord.id: {last_obs.id}' + ) raise Exception(f'Unable to submit next cadenced observation: {form.errors}') # Creation of corresponding ObservationRecord objects for the observations @@ -101,12 +120,11 @@ def run(self): ) # Add ObservationRecords to the DynamicCadence self.dynamic_cadence.observation_group.observation_records.add(record) - self.dynamic_cadence.observation_group.save() new_observations.append(record) + self.dynamic_cadence.observation_group.save() # Update the status of the ObservationRecords in the DB for obsr in new_observations: - facility = get_service_class(obsr.facility)() facility.update_observation_status(obsr.observation_id) obsr.refresh_from_db() # commit the updated observation status @@ -114,25 +132,25 @@ def run(self): def advance_window(self, observation_payload, start_keyword='start', end_keyword='end'): cadence_frequency = self.dynamic_cadence.cadence_parameters.get('cadence_frequency') - if not cadence_frequency: + if cadence_frequency is None: raise Exception(f'The {self.name} strategy requires a cadence_frequency cadence_parameter.') advance_window_hours = cadence_frequency + window_min = getattr(settings, 'OBS_WINDOW_MINIMUM', 24) + window_length = min(cadence_frequency, window_min) - # Window length for the observation should be every 24 hours unless the frequency is less than 24 - # then just the cadence frequency (24 hour window defined in the settings) - if settings.OBS_WINDOW_MINIMUM: - window_length = settings.OBS_WINDOW_MINIMUM - if window_length > cadence_frequency: - window_length = cadence_frequency - else: - window_length = cadence_frequency + scheduled_end = observation_payload['scheduled_end'] + + if isinstance(scheduled_end, str): + scheduled_end = parse(scheduled_end) - # Define new start to be at the end of the previous observation + cadence in hours - new_start = parse(observation_payload['scheduled_end']) + timedelta(hours=advance_window_hours) - if new_start < datetime.now(): # Ensure that the new window isn't in the past - new_start = datetime.now() - new_end = new_start + window_length + if timezone.is_naive(scheduled_end): + scheduled_end = timezone.make_aware(scheduled_end) + + new_start = scheduled_end + timedelta(hours=advance_window_hours) + if new_start < timezone.now(): # Ensure that the new window isn't in the past + new_start = timezone.now() + new_end = new_start + timedelta(hours=window_length) observation_payload[start_keyword] = new_start.isoformat() observation_payload[end_keyword] = new_end.isoformat() - return observation_payload + return observation_payload \ No newline at end of file diff --git a/tom_observations/cadences/retry_failed_observations.py b/tom_observations/cadences/retry_failed_observations.py index 001ed780d..620ee6c05 100644 --- a/tom_observations/cadences/retry_failed_observations.py +++ b/tom_observations/cadences/retry_failed_observations.py @@ -2,6 +2,7 @@ from dateutil.parser import parse import logging from django.conf import settings +from django.utils import timezone from tom_observations.cadence import BaseCadenceForm, CadenceStrategy from tom_observations.models import ObservationRecord, DynamicCadence @@ -14,9 +15,9 @@ class RetryFailedObservationsForm(BaseCadenceForm): pass -class RetryFailedObservationsStrategy(CadenceStrategy): +class BaseRetryFailedObservationsStrategy(CadenceStrategy): """ - The RetryFailedObservationsStrategy immediately re-submits all observations within an observation group a certain + The BaseRetryFailedObservationsStrategy immediately re-submits all observations within an observation group a certain number of hours later, as specified by ``advance_window_hours``. This strategy requires the DynamicCadence to have a parameter ``cadence_frequency``. @@ -27,74 +28,169 @@ class RetryFailedObservationsStrategy(CadenceStrategy): form = RetryFailedObservationsForm def run(self): - records = self.dynamic_cadence.observation_group.observation_records.all().order_by('-created') - last_obs = records.first() + records = self.dynamic_cadence.observation_group.observation_records.all().order_by('created') + first_obs = records.first() + last_obs = records.last() - if not last_obs: + if not first_obs or not last_obs: return - facility_class = get_service_class(last_obs.facility) - facility = facility_class() - start_keyword, end_keyword = facility.get_start_end_keywords() + facility = get_service_class(last_obs.facility)() facility.update_observation_status(last_obs.observation_id) last_obs.refresh_from_db() - if not last_obs.terminal: #observation is still pending, do nothing + if not last_obs.terminal: + return + + if last_obs.status == 'COMPLETED': + self.dynamic_cadence.active = False + self.dynamic_cadence.save() + logger.info(f'Observation {last_obs} completed; turned off dynamic cadence') + return self.notify_success(last_obs) + + if last_obs.status == 'CANCELED': + self.dynamic_cadence.active = False + self.dynamic_cadence.save() + logger.info(f'Observation {last_obs} was canceled, stopping dynamic cadence') return - - elif not last_obs.failed: #observation succeeded + + if not self.retry_observation(first_obs, last_obs, facility): + self.dynamic_cadence.active = False + self.dynamic_cadence.save() + logger.info( + f'Stopping retry cadence for observation group ' + f'{self.dynamic_cadence.observation_group.id}' + ) + return + + return self.submit_retry_observation(first_obs, last_obs, facility) + + def submit_retry_observation(self, first_obs, last_obs, facility): + observation_payload = last_obs.parameters.copy() + + start_keyword, end_keyword = facility.get_start_end_keywords() + observation_payload = self.advance_window( + observation_payload, start_keyword=start_keyword, end_keyword=end_keyword, first_obs=first_obs, facility=facility + ) + + if observation_payload is None: self.dynamic_cadence.active = False self.dynamic_cadence.save() - return 'COMPLETED' + logger.info( + f'No retry window remaining for observation group ' + f'{self.dynamic_cadence.observation_group.id}; deactivated silently' + ) + return + + obs_type = observation_payload.get('observation_type') + form = facility.get_form(obs_type)(observation_payload) + + if not form.is_valid(): + logger.error( + msg=f'Unable to submit next observation: {form.errors} ' + f'for ObservationRecord.id: {last_obs.id}' + ) + raise Exception(f'Unable to submit next observation: {form.errors}') - else: #observation failed, submit a new one - observation_payload = last_obs.parameters.copy() + observation_ids = facility.submit_observation(form.observation_payload()) + new_observations = [] - observation_payload = self.advance_window( - observation_payload, start_keyword=start_keyword, end_keyword=end_keyword + for observation_id in observation_ids: + record = ObservationRecord.objects.create( + target=last_obs.target, + facility=facility.name, + parameters=observation_payload, + observation_id=observation_id, ) - - obs_type = observation_payload.get('observation_type') - form = facility.get_form(obs_type)(observation_payload) - - if not form.is_valid(): - logger.error(msg=f'Unable to submit next cadenced observation: {form.errors}') - raise Exception(f'Unable to submit next cadenced observation: {form.errors}') - - observation_ids = facility.submit_observation(form.observation_payload()) - new_observations = [] - - for observation_id in observation_ids: - record = ObservationRecord.objects.create( - target=last_obs.target, - facility=facility.name, - parameters=observation_payload, - observation_id=observation_id - ) - self.dynamic_cadence.observation_group.observation_records.add(record) - new_observations.append(record) - - self.dynamic_cadence.observation_group.save() - - for obsr in new_observations: - facility.update_observation_status(obsr.observation_id) - obsr.refresh_from_db() - - return new_observations - - def advance_window(self, observation_payload, start_keyword='start', end_keyword='end'): + self.dynamic_cadence.observation_group.observation_records.add(record) + new_observations.append(record) + + self.dynamic_cadence.observation_group.save() + + for obsr in new_observations: + facility.update_observation_status(obsr.observation_id) + obsr.refresh_from_db() + + return new_observations + + def advance_window(self, observation_payload, start_keyword='start', end_keyword='end', first_obs=None, facility=None): cadence_frequency = self.dynamic_cadence.cadence_parameters.get('cadence_frequency') - if not cadence_frequency: - raise Exception(f'The {self.name} strategy requires a cadence_frequency cadence_parameter.') - if settings.OBS_WINDOW_MINIMUM: - window_length = settings.OBS_WINDOW_MINIMUM - if window_length > cadence_frequency: - window_length = cadence_frequency - else: - window_length = cadence_frequency - new_start = parse(observation_payload[start_keyword]) + timedelta(hours=window_length) - new_end = parse(observation_payload[end_keyword]) + timedelta(hours=window_length) + if cadence_frequency is None: + raise Exception( + f'The {self.name} strategy requires a cadence_frequency cadence_parameter.' + ) + + window_min = getattr(settings, 'OBS_WINDOW_MINIMUM', 24) + window_length = min(cadence_frequency, window_min) + + new_start = timezone.now() + new_end = new_start + timedelta(hours=window_length) + observation_payload[start_keyword] = new_start.isoformat() observation_payload[end_keyword] = new_end.isoformat() + return observation_payload + + +class RetryFailedObservationsStrategy(BaseRetryFailedObservationsStrategy): + """ + Retry indefinitely until the observation succeeds. + """ + pass + + +class RetryUntilDeadlineStrategy(BaseRetryFailedObservationsStrategy): + """ + Retry in short windows until either the observation succeeds, or the + original cadence_frequency interval has elapsed. + """ + + def retry_observation(self, first_obs, last_obs, facility): + deadline = self.get_deadline(first_obs, facility) + return timezone.now() < deadline + + def advance_window(self, observation_payload, start_keyword='start', end_keyword='end', first_obs=None, facility=None): + cadence_frequency = self.dynamic_cadence.cadence_parameters.get('cadence_frequency') + if cadence_frequency is None: + raise Exception( + f'The {self.name} strategy requires a cadence_frequency cadence_parameter.' + ) + window_min = getattr(settings, 'OBS_WINDOW_MINIMUM', 24) + window_length = min(cadence_frequency, window_min) + + deadline = self.get_deadline(first_obs, facility) + new_start = timezone.now() + + if new_start >= deadline: + return None + + new_end = min(new_start + timedelta(hours=window_length), deadline) + + if new_end <= new_start: + return None + + observation_payload[start_keyword] = new_start.isoformat() + observation_payload[end_keyword] = new_end.isoformat() return observation_payload + + def get_deadline(self, first_obs, facility): + cadence_frequency = self.dynamic_cadence.cadence_parameters.get('cadence_frequency') + if cadence_frequency is None: + raise Exception( + f'The {self.name} strategy requires a cadence_frequency cadence_parameter.' + ) + + start_keyword, _ = facility.get_start_end_keywords() + start_value = first_obs.parameters.get(start_keyword) + + if not start_value: + raise Exception( + f'Could not determine original start time for ' + f'ObservationRecord.id={first_obs.id}' + ) + + original_start = parse(start_value) if isinstance(start_value, str) else start_value + if timezone.is_naive(original_start): + original_start = timezone.make_aware(original_start) + + return original_start + timedelta(hours=cadence_frequency) \ No newline at end of file diff --git a/tom_observations/facilities/lco.py b/tom_observations/facilities/lco.py index 4f4a6897d..34500ceac 100644 --- a/tom_observations/facilities/lco.py +++ b/tom_observations/facilities/lco.py @@ -789,7 +789,7 @@ class LCOPhotometricSequenceForm(LCOOldStyleObservationForm): The form is modeled after the Supernova Exchange application's Photometric Sequence Request Form, and allows the configuration of multiple filters, as well as a more intuitive proactive cadence form. """ - valid_instruments = ['1M0-SCICAM-SINISTRO', '0M4-SCICAM-SBIG', '2M0-SPECTRAL-AG'] + valid_instruments = ['1M0-SCICAM-SINISTRO', '0M4-SCICAM-SBIG', '2M0-SPECTRAL-AG', '0M4-SCICAM-QHY600'] valid_filters = ['U', 'B', 'V', 'R', 'I', 'up', 'gp', 'rp', 'ip', 'zs', 'w', 'unknown'] cadence_frequency = forms.IntegerField(required=True, help_text='in hours') @@ -808,7 +808,11 @@ def __init__(self, *args, **kwargs): # Massage cadence form to be SNEx-styled self.fields['cadence_strategy'] = forms.ChoiceField( - choices=[('', 'Once in the next'), ('ResumeCadenceAfterFailureStrategy', 'Repeating every')], + choices=[ + ('ResumeCadenceAfterFailureStrategy', 'Repeating every'), + ('RetryUntilDeadlineStrategy', 'Once in the next'), + ('RetryFailedObservationsStrategy', 'Retry until successful'), + ], required=False, ) for field_name in ['exposure_time', 'exposure_count', 'filter']: @@ -816,8 +820,7 @@ def __init__(self, *args, **kwargs): if self.fields.get('groups'): self.fields['groups'].label = 'Data granted to' for field_name in ['start', 'end']: - self.fields[field_name].widget = forms.HiddenInput() - self.fields[field_name].required = False + self.fields[field_name] = forms.CharField(required=False, widget=forms.HiddenInput()) self.helper.layout = Layout( Row( @@ -882,14 +885,11 @@ def clean(self): - Adds an end time that corresponds with the cadence frequency """ cleaned_data = super().clean() + logger.info(f'cleaned data: {cleaned_data}') start = cleaned_data.get('start') cadence_frequency = cleaned_data['cadence_frequency'] - if settings.OBS_WINDOW_MINIMUM: - window_length = settings.OBS_WINDOW_MINIMUM - if window_length > cadence_frequency: - window_length = cadence_frequency - else: - window_length = cadence_frequency + window_min = getattr(settings, 'OBS_WINDOW_MINIMUM', 24) + window_length = min(window_min, cadence_frequency) cleaned_data['end'] = datetime.strftime(parse(start) + timedelta(hours=window_length), '%Y-%m-%dT%H:%M:%S') @@ -974,7 +974,11 @@ def __init__(self, *args, **kwargs): self.fields['name'].widget.attrs['placeholder'] = 'Name' self.fields['min_lunar_distance'].widget.attrs['placeholder'] = 'Degrees' self.fields['cadence_strategy'] = forms.ChoiceField( - choices=[('', 'Once in the next'), ('ResumeCadenceAfterFailureStrategy', 'Repeating every')], + choices=[ + ('ResumeCadenceAfterFailureStrategy', 'Repeating every'), + ('RetryUntilDeadlineStrategy', 'Once in the next'), + ('RetryFailedObservationsStrategy', 'Retry until successful'), + ], required=False, label='' ) @@ -986,8 +990,7 @@ def __init__(self, *args, **kwargs): if self.fields.get('groups'): self.fields['groups'].label = 'Data granted to' for field_name in ['start', 'end']: - self.fields[field_name].widget = forms.HiddenInput() - self.fields[field_name].required = False + self.fields[field_name] = forms.CharField(required=False, widget=forms.HiddenInput()) self.helper.layout = Layout( Div( @@ -1068,12 +1071,10 @@ def clean(self): cleaned_data['instrument_type'] = '2M0-FLOYDS-SCICAM' # SNEx only submits spectra to FLOYDS start = cleaned_data.get('start') - if settings.OBS_WINDOW_MINIMUM: - min_window = settings.OBS_WINDOW_MINIMUM - else: - min_window = 24 - window_length = min_window if cleaned_data['cadence_frequency'] > min_window else cleaned_data['cadence_frequency'] - + cadence_frequency = cleaned_data['cadence_frequency'] + window_min = getattr(settings, 'OBS_WINDOW_MINIMUM', 24) + window_length = min(cadence_frequency, window_min) + cleaned_data['end'] = datetime.strftime(parse(start) + timedelta(hours=window_length), '%Y-%m-%dT%H:%M:%S') From a3a8080e8234955176eee06cba923c2c2dfadc5b Mon Sep 17 00:00:00 2001 From: Moira Andrews Date: Thu, 25 Jun 2026 10:36:29 -0700 Subject: [PATCH 18/24] fixed canceled cadence not turning off dc, updated tests --- .../cadences/resume_cadence_after_failure.py | 6 +- .../cadences/retry_failed_observations.py | 5 + tom_observations/tests/test_cadence.py | 93 ++++++++++++------- 3 files changed, 71 insertions(+), 33 deletions(-) diff --git a/tom_observations/cadences/resume_cadence_after_failure.py b/tom_observations/cadences/resume_cadence_after_failure.py index e480ada11..4028ac226 100644 --- a/tom_observations/cadences/resume_cadence_after_failure.py +++ b/tom_observations/cadences/resume_cadence_after_failure.py @@ -55,9 +55,11 @@ def run(self): # If the observation hasn't finished, do nothing if not last_obs.terminal: return - + if last_obs.status == 'CANCELED': - logger.info(f'Observation {last_obs} was canceled, not resuming cadence') + self.dynamic_cadence.active = False + self.dynamic_cadence.save() + logger.info(f'Observation {last_obs} was canceled, stopping dynamic cadence') return # Boilerplate to get necessary properties for future calls diff --git a/tom_observations/cadences/retry_failed_observations.py b/tom_observations/cadences/retry_failed_observations.py index 620ee6c05..be290c58f 100644 --- a/tom_observations/cadences/retry_failed_observations.py +++ b/tom_observations/cadences/retry_failed_observations.py @@ -27,6 +27,11 @@ class BaseRetryFailedObservationsStrategy(CadenceStrategy): cadence.""" form = RetryFailedObservationsForm + def notify_success(self): + ''' + Function to add a call to slack or email on successful single time observations + ''' + def run(self): records = self.dynamic_cadence.observation_group.observation_records.all().order_by('created') first_obs = records.first() diff --git a/tom_observations/tests/test_cadence.py b/tom_observations/tests/test_cadence.py index e6ff65482..3ecb2295f 100644 --- a/tom_observations/tests/test_cadence.py +++ b/tom_observations/tests/test_cadence.py @@ -1,4 +1,5 @@ from django.test import TestCase +from django.utils import timezone from unittest.mock import patch from datetime import datetime, timedelta from dateutil.parser import parse @@ -56,46 +57,57 @@ def setUp(self): self.group.observation_records.add(*observing_records) self.group.save() self.dynamic_cadence = DynamicCadence.objects.create( - cadence_strategy='Test Strategy', cadence_parameters={'cadence_frequency': 72}, active=True, - observation_group=self.group) + cadence_strategy='RetryFailedObservationsStrategy', cadence_parameters={'cadence_frequency': 72}, + active=True, observation_group=self.group) - @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', return_value={'state': 'CANCELED', - 'scheduled_start': None, 'scheduled_end': None}) - def test_retry_when_failed_cadence_failed_obs(self, patch1, patch2, patch3, patch4, mock_get_obs_status, mock_validate_obs): + @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', + return_value={'state': 'WINDOW_EXPIRED', 'scheduled_start': None, 'scheduled_end': None}) + def test_retry_when_failed_cadence_failed_obs(self, mock_get_obs_status, mock_validate_obs, mock_submit_obs, + mock_proposal_choices, mock_get_insts): mock_validate_obs.return_value = {} num_records = self.group.observation_records.count() - observing_record = self.group.observation_records.first() - observing_record.status = 'CANCELED' - observing_record.save() strategy = RetryFailedObservationsStrategy(self.dynamic_cadence) new_records = strategy.run() self.group.refresh_from_db() - # Make sure the candence run created a new observation. + self.dynamic_cadence.refresh_from_db() self.assertEqual(num_records + 1, self.group.observation_records.count()) - # assert that the newly added observation record has a window of exactly 3 days - # later than the canceled observation. - self.assertEqual( - parse(observing_record.parameters['start']), - parse(new_records[0].parameters['start']) - timedelta(days=3) + self.assertAlmostEqual( + parse(new_records[0].parameters['start']), + timezone.now(), + delta=timedelta(seconds=5) ) - - @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', return_value={'state': 'CANCELED', - 'scheduled_start': None, 'scheduled_end': None}) - def test_retry_when_failed_cadence_successful_obs(self, patch1, patch2, patch3, patch4, mock_get_obs_status, mock_validate_obs): + self.assertTrue(self.dynamic_cadence.active) + + @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', + return_value={'state': 'COMPLETED', 'scheduled_start': None, 'scheduled_end': None}) + def test_retry_when_failed_cadence_successful_obs(self, mock_get_obs_status, mock_validate_obs, mock_submit_obs, + mock_proposal_choices, mock_get_insts): mock_validate_obs.return_value = {} - observing_record = self.group.observation_records.first() - observing_record.status = 'COMPLETE' - observing_record.save() + num_records = self.group.observation_records.count() strategy = RetryFailedObservationsStrategy(self.dynamic_cadence) new_records = strategy.run() self.group.refresh_from_db() - # Make sure the candence returned 'COMPLETED' - self.assertEqual(new_records, 'COMPLETED') - # Make sure the dynamic cadence was turned off - self.assertEqual(self.dynamic_cadence.active, False) + self.dynamic_cadence.refresh_from_db() + self.assertIsNone(new_records) + self.assertEqual(num_records, self.group.observation_records.count()) + self.assertFalse(self.dynamic_cadence.active) + + @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', + return_value={'state': 'CANCELED', 'scheduled_start': None, 'scheduled_end': None}) + def test_retry_when_canceled_deactivates(self, mock_get_obs_status, mock_validate_obs, mock_submit_obs, + mock_proposal_choices, mock_get_insts): + mock_validate_obs.return_value = {} + num_records = self.group.observation_records.count() + strategy = RetryFailedObservationsStrategy(self.dynamic_cadence) + new_records = strategy.run() + self.group.refresh_from_db() + self.dynamic_cadence.refresh_from_db() + self.assertIsNone(new_records) + self.assertEqual(num_records, self.group.observation_records.count()) + self.assertFalse(self.dynamic_cadence.active) @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', return_value={'state': 'CANCELED', 'scheduled_start': None, 'scheduled_end': None}) @@ -107,12 +119,14 @@ def test_resume_when_failed_cadence_failed_obs(self, mock_get_obs_status, mock_v strategy = ResumeCadenceAfterFailureStrategy(self.dynamic_cadence) new_records = strategy.run() self.group.refresh_from_db() + self.dynamic_cadence.refresh_from_db() self.assertEqual(num_records + 1, self.group.observation_records.count()) self.assertAlmostEqual( - parse(new_records[0].parameters["start"]), - datetime.now(), + parse(new_records[0].parameters['start']), + timezone.now(), delta=timedelta(seconds=5), ) + self.assertTrue(self.dynamic_cadence.active) @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', return_value={'state': 'COMPLETED', 'scheduled_start': None, 'scheduled_end': None}) @@ -127,8 +141,9 @@ def test_resume_when_failed_cadence_successful_obs(self, mock_get_obs_status, mo self.group.refresh_from_db() self.assertEqual(num_records + 1, self.group.observation_records.count()) self.assertAlmostEqual( - parse(obsr.parameters['start']).replace(second=0, microsecond=0), - parse(new_records[0].parameters['start']).replace(second=0, microsecond=0) - timedelta(days=3) + parse(new_records[0].parameters['start']), + timezone.make_aware(parse(obsr.parameters['end'])) + timedelta(hours=72), + delta=timedelta(seconds=5) ) @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', return_value={'state': 'COMPLETED', @@ -147,10 +162,26 @@ def test_resume_when_failed_cadence_invalid_date(self, mock_get_obs_status, mock self.group.refresh_from_db() self.assertEqual(num_records + 1, self.group.observation_records.count()) self.assertAlmostEqual( - datetime.now().replace(second=0, microsecond=0), - parse(new_records[0].parameters['start']).replace(second=0, microsecond=0) + timezone.now(), + parse(new_records[0].parameters['start']), + delta=timedelta(seconds=5) ) + @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', + return_value={'state': 'CANCELED', 'scheduled_start': None, 'scheduled_end': None}) + def test_resume_when_canceled_deactivates(self, mock_get_obs_status, mock_validate_obs, mock_submit_obs, + mock_proposal_choices, mock_get_insts): + mock_validate_obs.return_value = {} + num_records = self.group.observation_records.count() + + strategy = ResumeCadenceAfterFailureStrategy(self.dynamic_cadence) + new_records = strategy.run() + self.group.refresh_from_db() + self.dynamic_cadence.refresh_from_db() + self.assertIsNone(new_records) + self.assertEqual(num_records, self.group.observation_records.count()) + self.assertFalse(self.dynamic_cadence.active) + @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', return_value={'state': 'COMPLETED', 'scheduled_start': None, 'scheduled_end': None}) def test_resume_when_failed_cadence_obs_invalid(self, mock_get_obs_status, mock_validate_obs, mock_submit_obs, From 46e25499464a2086d5a571c7d51854dd68165018 Mon Sep 17 00:00:00 2001 From: Moira Andrews Date: Thu, 25 Jun 2026 14:23:55 -0700 Subject: [PATCH 19/24] added cadence tests --- tom_observations/tests/test_cadence.py | 71 +++++++++++++++++++++++++- 1 file changed, 69 insertions(+), 2 deletions(-) diff --git a/tom_observations/tests/test_cadence.py b/tom_observations/tests/test_cadence.py index 3ecb2295f..7ac9c1397 100644 --- a/tom_observations/tests/test_cadence.py +++ b/tom_observations/tests/test_cadence.py @@ -7,7 +7,7 @@ from .factories import ObservingRecordFactory, SiderealTargetFactory from tom_observations.models import ObservationGroup, DynamicCadence from tom_observations.cadences.resume_cadence_after_failure import ResumeCadenceAfterFailureStrategy -from tom_observations.cadences.retry_failed_observations import RetryFailedObservationsStrategy +from tom_observations.cadences.retry_failed_observations import RetryFailedObservationsStrategy, RetryUntilDeadlineStrategy mock_instruments = { @@ -109,7 +109,7 @@ def test_retry_when_canceled_deactivates(self, mock_get_obs_status, mock_validat self.assertEqual(num_records, self.group.observation_records.count()) self.assertFalse(self.dynamic_cadence.active) - @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', return_value={'state': 'CANCELED', + @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', return_value={'state': 'WINDOW_EXPIRED', 'scheduled_start': None, 'scheduled_end': None}) def test_resume_when_failed_cadence_failed_obs(self, mock_get_obs_status, mock_validate_obs, mock_submit_obs, mock_proposal_choices, mock_get_insts): @@ -191,3 +191,70 @@ def test_resume_when_failed_cadence_obs_invalid(self, mock_get_obs_status, mock_ strategy = ResumeCadenceAfterFailureStrategy(self.dynamic_cadence) with self.assertRaises(Exception): strategy.run() + @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', + return_value={'state': 'PENDING', 'scheduled_start': None, 'scheduled_end': None}) + def test_retry_when_not_terminal_noop(self, mock_get_obs_status, mock_validate_obs, mock_submit_obs, + mock_proposal_choices, mock_get_insts): + mock_validate_obs.return_value = {} + num_records = self.group.observation_records.count() + + strategy = RetryFailedObservationsStrategy(self.dynamic_cadence) + new_records = strategy.run() + self.group.refresh_from_db() + self.dynamic_cadence.refresh_from_db() + self.assertIsNone(new_records) + self.assertEqual(num_records, self.group.observation_records.count()) + self.assertTrue(self.dynamic_cadence.active) + + @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', + return_value={'state': 'PENDING', 'scheduled_start': None, 'scheduled_end': None}) + def test_resume_when_not_terminal_noop(self, mock_get_obs_status, mock_validate_obs, mock_submit_obs, + mock_proposal_choices, mock_get_insts): + mock_validate_obs.return_value = {} + num_records = self.group.observation_records.count() + + strategy = ResumeCadenceAfterFailureStrategy(self.dynamic_cadence) + new_records = strategy.run() + self.group.refresh_from_db() + self.dynamic_cadence.refresh_from_db() + self.assertIsNone(new_records) + self.assertEqual(num_records, self.group.observation_records.count()) + self.assertTrue(self.dynamic_cadence.active) + + @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', + return_value={'state': 'WINDOW_EXPIRED', 'scheduled_start': None, 'scheduled_end': None}) + def test_retry_until_deadline_before_deadline_retries(self, mock_get_obs_status, mock_validate_obs, mock_submit_obs, + mock_proposal_choices, mock_get_insts): + mock_validate_obs.return_value = {} + num_records = self.group.observation_records.count() + + strategy = RetryUntilDeadlineStrategy(self.dynamic_cadence) + new_records = strategy.run() + self.group.refresh_from_db() + self.dynamic_cadence.refresh_from_db() + self.assertEqual(num_records + 1, self.group.observation_records.count()) + self.assertAlmostEqual( + parse(new_records[0].parameters['start']), + timezone.now(), + delta=timedelta(seconds=5) + ) + self.assertTrue(self.dynamic_cadence.active) + + @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', + return_value={'state': 'WINDOW_EXPIRED', 'scheduled_start': None, 'scheduled_end': None}) + def test_retry_until_deadline_after_deadline_deactivates(self, mock_get_obs_status, mock_validate_obs, + mock_submit_obs, mock_proposal_choices, mock_get_insts): + mock_validate_obs.return_value = {} + first_obs = self.group.observation_records.order_by('created').first() + first_obs.parameters = {**first_obs.parameters, + 'start': (datetime.now() - timedelta(hours=100)).strftime('%Y-%m-%dT%H:%M:%S')} + first_obs.save() + num_records = self.group.observation_records.count() + + strategy = RetryUntilDeadlineStrategy(self.dynamic_cadence) + new_records = strategy.run() + self.group.refresh_from_db() + self.dynamic_cadence.refresh_from_db() + self.assertIsNone(new_records) + self.assertEqual(num_records, self.group.observation_records.count()) + self.assertFalse(self.dynamic_cadence.active) \ No newline at end of file From 800e514dcfbcb0d70f96087e4be17140249df645 Mon Sep 17 00:00:00 2001 From: Moira Andrews Date: Thu, 25 Jun 2026 14:44:48 -0700 Subject: [PATCH 20/24] fixed some test bugs --- .../cadences/retry_failed_observations.py | 14 +++++++++----- tom_observations/tests/test_cadence.py | 2 +- 2 files changed, 10 insertions(+), 6 deletions(-) diff --git a/tom_observations/cadences/retry_failed_observations.py b/tom_observations/cadences/retry_failed_observations.py index be290c58f..af6c7830c 100644 --- a/tom_observations/cadences/retry_failed_observations.py +++ b/tom_observations/cadences/retry_failed_observations.py @@ -27,10 +27,17 @@ class BaseRetryFailedObservationsStrategy(CadenceStrategy): cadence.""" form = RetryFailedObservationsForm - def notify_success(self): + def retry_observation(self, first_obs, last_obs, facility): + ''' + Default retry observations, for BaseRetry strategy (retry until successful), the default is to always retry the observations until obtained + ''' + return True + + def notify_success(self, obs): ''' Function to add a call to slack or email on successful single time observations ''' + return def run(self): records = self.dynamic_cadence.observation_group.observation_records.all().order_by('created') @@ -62,10 +69,7 @@ def run(self): if not self.retry_observation(first_obs, last_obs, facility): self.dynamic_cadence.active = False self.dynamic_cadence.save() - logger.info( - f'Stopping retry cadence for observation group ' - f'{self.dynamic_cadence.observation_group.id}' - ) + logger.info(f'Stopping retry cadence for observation group {self.dynamic_cadence.observation_group.id}') return return self.submit_retry_observation(first_obs, last_obs, facility) diff --git a/tom_observations/tests/test_cadence.py b/tom_observations/tests/test_cadence.py index 7ac9c1397..0b75f9aa2 100644 --- a/tom_observations/tests/test_cadence.py +++ b/tom_observations/tests/test_cadence.py @@ -257,4 +257,4 @@ def test_retry_until_deadline_after_deadline_deactivates(self, mock_get_obs_stat self.dynamic_cadence.refresh_from_db() self.assertIsNone(new_records) self.assertEqual(num_records, self.group.observation_records.count()) - self.assertFalse(self.dynamic_cadence.active) \ No newline at end of file + self.assertFalse(self.dynamic_cadence.active) From 135cd5d31f4adab4361a6f9151e6872649114bb4 Mon Sep 17 00:00:00 2001 From: Moira Andrews Date: Thu, 25 Jun 2026 14:47:02 -0700 Subject: [PATCH 21/24] linting fix --- tom_observations/tests/test_cadence.py | 1 + 1 file changed, 1 insertion(+) diff --git a/tom_observations/tests/test_cadence.py b/tom_observations/tests/test_cadence.py index 0b75f9aa2..3fc41acda 100644 --- a/tom_observations/tests/test_cadence.py +++ b/tom_observations/tests/test_cadence.py @@ -191,6 +191,7 @@ def test_resume_when_failed_cadence_obs_invalid(self, mock_get_obs_status, mock_ strategy = ResumeCadenceAfterFailureStrategy(self.dynamic_cadence) with self.assertRaises(Exception): strategy.run() + @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', return_value={'state': 'PENDING', 'scheduled_start': None, 'scheduled_end': None}) def test_retry_when_not_terminal_noop(self, mock_get_obs_status, mock_validate_obs, mock_submit_obs, From bd29b2c46dd3ed45e897e41a408a522afa5d45d8 Mon Sep 17 00:00:00 2001 From: Moira Andrews Date: Thu, 25 Jun 2026 15:02:53 -0700 Subject: [PATCH 22/24] fixed linting --- .../cadences/resume_cadence_after_failure.py | 8 +++---- .../cadences/retry_failed_observations.py | 22 +++++++++++-------- tom_observations/facilities/lco.py | 2 +- tom_observations/tests/test_cadence.py | 14 +++++++----- 4 files changed, 27 insertions(+), 19 deletions(-) diff --git a/tom_observations/cadences/resume_cadence_after_failure.py b/tom_observations/cadences/resume_cadence_after_failure.py index 4028ac226..9b8536c44 100644 --- a/tom_observations/cadences/resume_cadence_after_failure.py +++ b/tom_observations/cadences/resume_cadence_after_failure.py @@ -1,4 +1,4 @@ -from datetime import datetime, timedelta +from datetime import timedelta from dateutil.parser import parse from django.conf import settings import logging @@ -55,7 +55,7 @@ def run(self): # If the observation hasn't finished, do nothing if not last_obs.terminal: return - + if last_obs.status == 'CANCELED': self.dynamic_cadence.active = False self.dynamic_cadence.save() @@ -128,7 +128,7 @@ def run(self): # Update the status of the ObservationRecords in the DB for obsr in new_observations: facility.update_observation_status(obsr.observation_id) - obsr.refresh_from_db() # commit the updated observation status + obsr.refresh_from_db() # commit the updated observation status return new_observations @@ -155,4 +155,4 @@ def advance_window(self, observation_payload, start_keyword='start', end_keyword observation_payload[start_keyword] = new_start.isoformat() observation_payload[end_keyword] = new_end.isoformat() - return observation_payload \ No newline at end of file + return observation_payload diff --git a/tom_observations/cadences/retry_failed_observations.py b/tom_observations/cadences/retry_failed_observations.py index af6c7830c..c4ee1110a 100644 --- a/tom_observations/cadences/retry_failed_observations.py +++ b/tom_observations/cadences/retry_failed_observations.py @@ -5,7 +5,7 @@ from django.utils import timezone from tom_observations.cadence import BaseCadenceForm, CadenceStrategy -from tom_observations.models import ObservationRecord, DynamicCadence +from tom_observations.models import ObservationRecord from tom_observations.facility import get_service_class logger = logging.getLogger(__name__) @@ -17,19 +17,20 @@ class RetryFailedObservationsForm(BaseCadenceForm): class BaseRetryFailedObservationsStrategy(CadenceStrategy): """ - The BaseRetryFailedObservationsStrategy immediately re-submits all observations within an observation group a certain - number of hours later, as specified by ``advance_window_hours``. + The BaseRetryFailedObservationsStrategy immediately re-submits all observations within an observation + group a certain number of hours later, as specified by ``advance_window_hours``. This strategy requires the DynamicCadence to have a parameter ``cadence_frequency``. """ name = 'Retry Failed Observations' - description = """This strategy immediately re-submits a cadenced observation without amending any other part of the - cadence.""" + description = """This strategy immediately re-submits a cadenced observation without amending any + other part of the cadence.""" form = RetryFailedObservationsForm def retry_observation(self, first_obs, last_obs, facility): ''' - Default retry observations, for BaseRetry strategy (retry until successful), the default is to always retry the observations until obtained + Default retry observations, for BaseRetry strategy (retry until successful), the default is to always + retry the observations until obtained ''' return True @@ -79,7 +80,8 @@ def submit_retry_observation(self, first_obs, last_obs, facility): start_keyword, end_keyword = facility.get_start_end_keywords() observation_payload = self.advance_window( - observation_payload, start_keyword=start_keyword, end_keyword=end_keyword, first_obs=first_obs, facility=facility + observation_payload, start_keyword=start_keyword, end_keyword=end_keyword, + first_obs=first_obs, facility=facility ) if observation_payload is None: @@ -122,7 +124,8 @@ def submit_retry_observation(self, first_obs, last_obs, facility): return new_observations - def advance_window(self, observation_payload, start_keyword='start', end_keyword='end', first_obs=None, facility=None): + def advance_window(self, observation_payload, + start_keyword='start', end_keyword='end', first_obs=None, facility=None): cadence_frequency = self.dynamic_cadence.cadence_parameters.get('cadence_frequency') if cadence_frequency is None: raise Exception( @@ -157,7 +160,8 @@ def retry_observation(self, first_obs, last_obs, facility): deadline = self.get_deadline(first_obs, facility) return timezone.now() < deadline - def advance_window(self, observation_payload, start_keyword='start', end_keyword='end', first_obs=None, facility=None): + def advance_window(self, observation_payload, + start_keyword='start', end_keyword='end', first_obs=None, facility=None): cadence_frequency = self.dynamic_cadence.cadence_parameters.get('cadence_frequency') if cadence_frequency is None: raise Exception( diff --git a/tom_observations/facilities/lco.py b/tom_observations/facilities/lco.py index 34500ceac..0daf6abfc 100644 --- a/tom_observations/facilities/lco.py +++ b/tom_observations/facilities/lco.py @@ -1074,7 +1074,7 @@ def clean(self): cadence_frequency = cleaned_data['cadence_frequency'] window_min = getattr(settings, 'OBS_WINDOW_MINIMUM', 24) window_length = min(cadence_frequency, window_min) - + cleaned_data['end'] = datetime.strftime(parse(start) + timedelta(hours=window_length), '%Y-%m-%dT%H:%M:%S') diff --git a/tom_observations/tests/test_cadence.py b/tom_observations/tests/test_cadence.py index 3fc41acda..6384f15fc 100644 --- a/tom_observations/tests/test_cadence.py +++ b/tom_observations/tests/test_cadence.py @@ -7,7 +7,9 @@ from .factories import ObservingRecordFactory, SiderealTargetFactory from tom_observations.models import ObservationGroup, DynamicCadence from tom_observations.cadences.resume_cadence_after_failure import ResumeCadenceAfterFailureStrategy -from tom_observations.cadences.retry_failed_observations import RetryFailedObservationsStrategy, RetryUntilDeadlineStrategy +from tom_observations.cadences.retry_failed_observations import ( + RetryFailedObservationsStrategy, RetryUntilDeadlineStrategy +) mock_instruments = { @@ -109,8 +111,8 @@ def test_retry_when_canceled_deactivates(self, mock_get_obs_status, mock_validat self.assertEqual(num_records, self.group.observation_records.count()) self.assertFalse(self.dynamic_cadence.active) - @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', return_value={'state': 'WINDOW_EXPIRED', - 'scheduled_start': None, 'scheduled_end': None}) + @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', + return_value={'state': 'WINDOW_EXPIRED', 'scheduled_start': None, 'scheduled_end': None}) def test_resume_when_failed_cadence_failed_obs(self, mock_get_obs_status, mock_validate_obs, mock_submit_obs, mock_proposal_choices, mock_get_insts): mock_validate_obs.return_value = {} @@ -224,7 +226,8 @@ def test_resume_when_not_terminal_noop(self, mock_get_obs_status, mock_validate_ @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', return_value={'state': 'WINDOW_EXPIRED', 'scheduled_start': None, 'scheduled_end': None}) - def test_retry_until_deadline_before_deadline_retries(self, mock_get_obs_status, mock_validate_obs, mock_submit_obs, + def test_retry_until_deadline_before_deadline_retries(self, mock_get_obs_status, + mock_validate_obs, mock_submit_obs, mock_proposal_choices, mock_get_insts): mock_validate_obs.return_value = {} num_records = self.group.observation_records.count() @@ -244,7 +247,8 @@ def test_retry_until_deadline_before_deadline_retries(self, mock_get_obs_status, @patch('tom_observations.facilities.lco.LCOFacility.get_observation_status', return_value={'state': 'WINDOW_EXPIRED', 'scheduled_start': None, 'scheduled_end': None}) def test_retry_until_deadline_after_deadline_deactivates(self, mock_get_obs_status, mock_validate_obs, - mock_submit_obs, mock_proposal_choices, mock_get_insts): + mock_submit_obs, mock_proposal_choices, + mock_get_insts): mock_validate_obs.return_value = {} first_obs = self.group.observation_records.order_by('created').first() first_obs.parameters = {**first_obs.parameters, From 9c4187c5ca3bbc0398a50c9cfe36118f77658196 Mon Sep 17 00:00:00 2001 From: Moira Andrews Date: Thu, 25 Jun 2026 15:05:45 -0700 Subject: [PATCH 23/24] fixed data bug --- tom_observations/cadences/retry_failed_observations.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tom_observations/cadences/retry_failed_observations.py b/tom_observations/cadences/retry_failed_observations.py index c4ee1110a..611e1f3aa 100644 --- a/tom_observations/cadences/retry_failed_observations.py +++ b/tom_observations/cadences/retry_failed_observations.py @@ -94,7 +94,7 @@ def submit_retry_observation(self, first_obs, last_obs, facility): return obs_type = observation_payload.get('observation_type') - form = facility.get_form(obs_type)(observation_payload) + form = facility.get_form(obs_type)(data=observation_payload) if not form.is_valid(): logger.error( From ba1d188c0391c2221c285430b4539d97497baf31 Mon Sep 17 00:00:00 2001 From: Moira Andrews Date: Thu, 25 Jun 2026 15:07:16 -0700 Subject: [PATCH 24/24] added extra line --- tom_observations/cadences/retry_failed_observations.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tom_observations/cadences/retry_failed_observations.py b/tom_observations/cadences/retry_failed_observations.py index 611e1f3aa..93b26ba43 100644 --- a/tom_observations/cadences/retry_failed_observations.py +++ b/tom_observations/cadences/retry_failed_observations.py @@ -206,4 +206,4 @@ def get_deadline(self, first_obs, facility): if timezone.is_naive(original_start): original_start = timezone.make_aware(original_start) - return original_start + timedelta(hours=cadence_frequency) \ No newline at end of file + return original_start + timedelta(hours=cadence_frequency)