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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
43 changes: 32 additions & 11 deletions imap_processing/cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
import sys
from abc import ABC, abstractmethod
from pathlib import Path
from time import sleep
from typing import final

import imap_data_access
Expand Down Expand Up @@ -472,6 +473,8 @@ def upload_products(self, products: list[Path]) -> None:
"""
Upload data products to the IMAP SDC.

If a 503 SlowDown error is reported from the Upload API, a retry will be made

Parameters
----------
products : list[Path]
Expand All @@ -483,18 +486,36 @@ def upload_products(self, products: list[Path]) -> None:
return

for filename in products:
try:
logger.info(f"Uploading file: {filename}")
imap_data_access.upload(filename)
except IMAPDataAccessError as e:
message = str(e)
if "FileAlreadyExists" in message and "409" in message:
logger.warning("Skipping upload of existing file, %s", filename)
continue
else:
max_retries = 3

for attempt in range(max_retries):
try:
logger.info(f"Uploading file: {filename}")
imap_data_access.upload(filename)
break

except IMAPDataAccessError as e:
message = str(e)

if "FileAlreadyExists" in message and "409" in message:
logger.warning(
"Skipping upload of existing file, %s", filename
)
break
elif "503" in message and "SlowDown" in message:
if attempt < max_retries - 1:
logger.warning(
"Upload busy. Waiting 5 seconds before retrying..."
)
sleep(5)
continue

logger.error(f"Upload failed with error: {message}")
except Exception as e:
logger.error(f"Upload failed unknown error: {e}")
raise

except Exception as e:
logger.error(f"Upload failed unknown error: {e}")
raise

@final
def process(self) -> None:
Expand Down
109 changes: 109 additions & 0 deletions imap_processing/tests/test_cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -937,3 +937,112 @@ def test_post_processing(
"naif0012.tls",
"imap_sclk_0001.tsc",
]


@mock.patch("imap_processing.cli.sleep")
@mock.patch("imap_processing.cli.filter_day_boundary_data")
@mock.patch("imap_processing.cli.swe_l1a")
def test_post_processing_upload_503_error(
mock_swe_l1a,
mock_filter,
mock_sleep,
mock_instrument_dependencies,
):
"""Test coverage for post processing when the upload fails with 503 error"""

mocks = mock_instrument_dependencies
mocks["mock_download"].return_value = "dependency0"
mocks["mock_write_cdf"].side_effect = [
"/path/to/imap_swe_l1a_test_20100105_v001.cdf"
]
mocks[
"mock_write_cdf"
].return_value = "/path/to/imap_swe_l1a_test_20100105_v001.cdf"
mocks["mock_query"].return_value = []

# Mocks a 503 error received from the upload API
mocks["mock_upload"].side_effect = imap_data_access.io.IMAPDataAccessError(
"503 Service Unavailable: "
"<title>503 Slow Down</title>"
"Code: SlowDown"
"Message: Please reduce your request rate."
)

test_ds = xr.Dataset()
mock_swe_l1a.return_value = [test_ds]
mock_filter.side_effect = lambda ds, _: ds
input_collection = ProcessingInputCollection(
ScienceInput("imap_swe_l0_raw_20100105_v001.pkts"),
SPICEInput("naif0012.tls", "imap_sclk_0001.tsc"),
)
mocks["mock_pre_processing"].return_value = input_collection

dependency_str = (
'[{"type": "science","files": ["imap_swe_l0_raw_20100105_v001.pkts"]}, '
'{"type": "spice", "files": ["naif0012.tls", "imap_sclk_0001.tsc"]}]'
)
instrument = Swe("l1a", "raw", dependency_str, "20100105", None, "v001", True)

# Checks that the upload failed and logs an error and raises an exception
with mock.patch("logging.Logger.error") as mock_error:
with pytest.raises(imap_data_access.io.IMAPDataAccessError):
instrument.process()

# Upload should attempt 3 times
assert mocks["mock_upload"].call_count == 3

# Sleep should be called 2 times after first two failures
assert mock_sleep.call_count == 2

# Checks the upload failure was logged
assert any(
"Upload failed with error" in str(call) for call in mock_error.call_args_list
)


@mock.patch("imap_processing.cli.filter_day_boundary_data")
@mock.patch("imap_processing.cli.swe_l1a")
def test_post_processing_upload_unknown_error(
mock_swe_l1a,
mock_filter,
mock_instrument_dependencies,
):
"""Test coverage for post processing when the upload fails with unknown error"""

mocks = mock_instrument_dependencies
mocks["mock_download"].return_value = "dependency0"
mocks["mock_write_cdf"].side_effect = [
"/path/to/imap_swe_l1a_test_20100105_v001.cdf"
]
mocks[
"mock_write_cdf"
].return_value = "/path/to/imap_swe_l1a_test_20100105_v001.cdf"
mocks["mock_query"].return_value = []

# Mocks an unknown error received from the upload API
mocks["mock_upload"].side_effect = RuntimeError("Unexpected failure")

test_ds = xr.Dataset()
mock_swe_l1a.return_value = [test_ds]
mock_filter.side_effect = lambda ds, _: ds
input_collection = ProcessingInputCollection(
ScienceInput("imap_swe_l0_raw_20100105_v001.pkts"),
SPICEInput("naif0012.tls", "imap_sclk_0001.tsc"),
)
mocks["mock_pre_processing"].return_value = input_collection

dependency_str = (
'[{"type": "science","files": ["imap_swe_l0_raw_20100105_v001.pkts"]}, '
'{"type": "spice", "files": ["naif0012.tls", "imap_sclk_0001.tsc"]}]'
)
instrument = Swe("l1a", "raw", dependency_str, "20100105", None, "v001", True)

# Checks that the upload failed and logs an error and raises an exception
with mock.patch("logging.Logger.error") as mock_error:
with pytest.raises(RuntimeError):
instrument.process()

# Checks the upload failure was logged
assert any(
"Upload failed unknown error" in str(call) for call in mock_error.call_args_list
)
Loading