"""Product scheduled update Logic CRUD operation.""" import time from datetime import datetime, timedelta from dateutil import parser from flask import g from oto import response from oto import status as http_status from sentry_sdk import capture_exception from timed_release.config import ( EXECUTE_STEP_FUNCTION_ARN, NOTIFY_LABEL_MANAGER_LAMBDA_ARN, SCHEDULER_GROUP, SCHEDULER_ROLE_ARN, TRIGGER_PRODUCT_UPDATE_LAMBDA_ARN, TRIGGER_SFN_WS_UPDATE_METADATA_ARN, WS_SCHEDULER_GROUP, WS_SCHEDULER_ROLE_ARN, WS_SEND_NOTIFICATION_LAMBDA_ARN, ) from timed_release.connectors import scheduler, sfn, sql from timed_release.constants import error, features from timed_release.constants.timed_release import ( OA, UPSERT_SCHEDULE, WARNING_EMAIL, ) from timed_release.models import scheduled_update from timed_release.utils import feature_flag as feature_flag_util @sql.db_session_wrap def upsert(product_id, data, orchard_identity_id, session=None): """Check if scheduled update data exists for the product and store. If data exists, update the data. If data does not exist, create new records. Handle scheduling logic when the timed releases workstation feature flag is enabled. This block performs the following steps: 1. Creates or updates the WS event schedule for the given product. 2. Deletes any existing OA schedule to avoid duplicate triggers. 3. Creates WS warning schedule that runs 2 hours before the main schedule. 4. Cleans up OA warning schedules after successful WS scheduling. Control flags: is_upsert_ws_event_schedule_called (bool): Set True after successfully creating the WS event schedule. Used in the exception handler to roll back WS schedules if a later step fails. is_delete_oa_event_schedule_called (bool): Set True after successfully deleting the OA schedule. Used in the exception handler to re-create OA schedule if WS scheduling ultimately fails. is_upsert_ws_warning_schedule_called (bool): Set True after success creating the WS warning schedule. Used to clean up warning schedules on failure. Exception handling: - If any step in WS scheduling fails, this block: • Rolls back the DB session. • Deletes any WS schedules that were created. • Optionally re-creates the OA schedule if it was deleted. • Returns an INTERNAL_ERROR response to the caller. Raises: Exception: Re-raises any exception from warning schedule creation to signal upstream failure and trigger cleanup logic. Args: product_id (int): Product id to update timed release data. data (dict): Schedule Update data to create/update with. orchard_identity_id (string): Users orchard identity id. Return: response.Response: timed release data on successful update or error response. """ sales_date_time_str = data['schedule_datetime'] sales_date_time = parser.isoparse(sales_date_time_str) time_offset = datetime.strptime(data['schedule_time_offset'], '%H:%M') process_at = sales_date_time - timedelta( hours=time_offset.hour, minutes=time_offset.minute) data['process_at_datetime'] = process_at.strftime('%Y-%m-%dT%H:%M:%SZ') warn_at = process_at - timedelta(hours=2) warn_at_datetime = warn_at.strftime('%Y-%m-%dT%H:%M:%SZ') result = scheduled_update.upsert( product_id, data, orchard_identity_id, session=session) is_timed_releases_ws_enabled = feature_flag_util.is_feature_enabled( features.CCM_TIMED_RELEASES_WORKSTATION, g.request_context ) is_upsert_ws_event_schedule_called = False is_delete_oa_event_schedule_called = False is_upsert_ws_warning_schedule_called = False if is_timed_releases_ws_enabled: try: scheduler.upsert_event_schedule( schedule_name=product_id, scheduler_group=WS_SCHEDULER_GROUP, scheduler_role_arn=WS_SCHEDULER_ROLE_ARN, schedule_datetime=data['process_at_datetime'], description=f'Metadata update for Product: {product_id}', target_arn=TRIGGER_SFN_WS_UPDATE_METADATA_ARN, payload={ 'productId': product_id, 'eventType': UPSERT_SCHEDULE, 'salesDateTime': sales_date_time_str, 'source': OA } ) is_upsert_ws_event_schedule_called = True scheduler.delete_event_schedule(product_id, SCHEDULER_GROUP) is_delete_oa_event_schedule_called = True try: warning_schedule_name = f'{product_id}-warning' scheduler.upsert_event_schedule( schedule_name=warning_schedule_name, scheduler_group=WS_SCHEDULER_GROUP, scheduler_role_arn=WS_SCHEDULER_ROLE_ARN, description=( f'Pre-metadata update warning for ' f'Product: {product_id}' ), schedule_datetime=warn_at_datetime, target_arn=WS_SEND_NOTIFICATION_LAMBDA_ARN, payload={ 'productId': product_id, 'eventProcessOn': data['process_at_datetime'], 'eventType': WARNING_EMAIL, 'salesDateTime': sales_date_time_str, 'source': OA } ) is_upsert_ws_warning_schedule_called = True scheduler.delete_event_schedule( warning_schedule_name, SCHEDULER_GROUP) except Exception as e: if is_upsert_ws_warning_schedule_called: scheduler.delete_event_schedule( warning_schedule_name, WS_SCHEDULER_GROUP) raise e except Exception: session.rollback() capture_exception() if is_upsert_ws_event_schedule_called: scheduler.delete_event_schedule(product_id, WS_SCHEDULER_GROUP) if is_delete_oa_event_schedule_called: scheduler.upsert_event_schedule( schedule_name=product_id, scheduler_group=SCHEDULER_GROUP, scheduler_role_arn=SCHEDULER_ROLE_ARN, description=f'Metadata update for Product: {product_id}', schedule_datetime=data['process_at_datetime'], target_arn=TRIGGER_PRODUCT_UPDATE_LAMBDA_ARN, payload={'productId': product_id} ) return response.create_error_response( code=error.INTERNAL_ERROR, message=error.ERROR_MESSAGE_UPSERTING_OA_SCHEDULE_FAILED, status=http_status.INTERNAL_ERROR ) else: try: scheduler.upsert_event_schedule( schedule_name=product_id, scheduler_group=SCHEDULER_GROUP, scheduler_role_arn=SCHEDULER_ROLE_ARN, description=f'Metadata update for Product: {product_id}', schedule_datetime=data['process_at_datetime'], target_arn=TRIGGER_PRODUCT_UPDATE_LAMBDA_ARN, payload={'productId': product_id} ) try: scheduler.upsert_event_schedule( schedule_name=f'{product_id}-warning', scheduler_group=SCHEDULER_GROUP, scheduler_role_arn=SCHEDULER_ROLE_ARN, description=( f'Pre-metadata update warning for ' f'Product: {product_id}' ), schedule_datetime=warn_at_datetime, target_arn=NOTIFY_LABEL_MANAGER_LAMBDA_ARN, payload={ 'productId': product_id, 'scheduledByIdentityId': orchard_identity_id, 'processAt': data['process_at_datetime'] } ) except Exception as e: scheduler.delete_event_schedule(product_id, SCHEDULER_GROUP) raise e except Exception: session.rollback() capture_exception() return response.create_error_response( code=error.INTERNAL_ERROR, message='error creating/updating event schedule', status=500 ) return result @sql.db_session_wrap def delete(product_id, session=None): """Delete a scheduled update by product id. Also delete the schedules from aws event bridge. Args: product_id (int): Product id to fetch schedule update data. Return: response.Response: successful delete or error response. """ result = scheduled_update.delete(product_id, session=session) is_timed_releases_ws_enabled = feature_flag_util.is_feature_enabled( features.CCM_TIMED_RELEASES_WORKSTATION, g.request_context ) try: if is_timed_releases_ws_enabled: scheduler.delete_event_schedule(product_id, WS_SCHEDULER_GROUP) try: scheduler.delete_event_schedule( f'{product_id}-warning', WS_SCHEDULER_GROUP) except Exception as e: schedule = scheduled_update.get(product_id) scheduler.upsert_event_schedule( schedule_name=product_id, scheduler_group=WS_SCHEDULER_GROUP, scheduler_role_arn=WS_SCHEDULER_ROLE_ARN, schedule_datetime=schedule['process_at_datetime'], description=f'Metadata update for Product: {product_id}', target_arn=TRIGGER_SFN_WS_UPDATE_METADATA_ARN, payload={ 'productId': product_id, 'eventType': UPSERT_SCHEDULE, 'salesDateTime': schedule['schedule_datetime'], 'source': OA } ) raise e else: scheduler.delete_event_schedule(product_id, SCHEDULER_GROUP) try: scheduler.delete_event_schedule( f'{product_id}-warning', SCHEDULER_GROUP) except Exception as e: schedule = scheduled_update.get(product_id) scheduler.upsert_event_schedule( schedule_name=product_id, scheduler_group=SCHEDULER_GROUP, scheduler_role_arn=SCHEDULER_ROLE_ARN, description=f'Metadata update for Product: {product_id}', schedule_datetime=schedule['process_at_datetime'], target_arn=TRIGGER_PRODUCT_UPDATE_LAMBDA_ARN, payload={'productId': product_id} ) raise e except Exception: session.rollback() capture_exception() return response.create_error_response( code=error.INTERNAL_ERROR, message='error deleting event schedule', status=500 ) return result def get_scheduled_update(product_id): """Retrieve scheduled update data for the product and store. Args: product_id (int): Product id to fetch schedule update data. Return: response.Response: schedule release data on successful update or error response. """ return scheduled_update.get(product_id) def execute_scheduled_update(product_id, name='None'): """Execute the scheduled update. Args: product_id (int): Product id to execute scheduled update for. Return: Response: response containing ResponseMetadata from execution """ get_schedule = scheduled_update.get(product_id) if not get_schedule: return get_schedule update_info = get_schedule.message payload = { 'productId': product_id, 'deliveryStoreIds': update_info.get('delivery_store_ids'), 'scheduledByIdentityId': update_info.get('scheduled_by'), 'updateBody': update_info.get('update') } timestamp = str(int(time.time())) suffix = f'-{name}' if name else '' result = sfn.execute_step_function( arn=EXECUTE_STEP_FUNCTION_ARN, name=f'productid-{product_id}-{timestamp}{suffix}', **payload ) return response.Response(message=result['ResponseMetadata'], status=200)