fix: report failed transfers instead of ending in success
Upload and download failures were printed and then forgotten: a batch that lost files still exited 0, and a file the client could not write was listed as if it had arrived. Nothing downstream could tell. - Return a verdict from every transfer worker and raise once at the end, so a run that lost a file exits 2. - Raise the error when a downloaded file cannot be written locally, instead of printing it and reporting the path as a success. - Write a download beside its destination and move it into place once complete, and refuse a destination that cannot be written, so a failed transfer no longer leaves a truncated file behind. - Carry on through the remaining sub-folders when a recursive download loses a file or cannot create a folder locally. - Name the file and give the reason in every failure message. - Count a failure the API layer did not raise, such as a file that disappeared between the directory walk and its turn to be sent. - Extract the directory walk from `Uploader.upload` to keep it within the complexity limit. - Cover all of the above in `tests/test_transfer.py`.
This commit is contained in:
+40
-4
@@ -1,5 +1,6 @@
|
||||
import mimetypes
|
||||
import os
|
||||
import threading
|
||||
from typing import Any, Final
|
||||
from unicodedata import normalize
|
||||
|
||||
@@ -146,18 +147,53 @@ class FilesApi(BaseApi):
|
||||
# print(self.__class__.__name__ + "::" + sys._getframe().f_code.co_name)
|
||||
url = file.download_url
|
||||
token_check(self.connection)
|
||||
# Refused before anything is fetched. The finished file is moved into place, and a
|
||||
# rename would replace a destination whose mode says it is protected.
|
||||
if os.path.exists(path):
|
||||
try:
|
||||
with open(path, "r+b"):
|
||||
pass
|
||||
except OSError as e:
|
||||
raise UnexpectedException(f"Cannot write `{path}`: {e}")
|
||||
response = self.connection.get(url, stream=True)
|
||||
self._raise_response_error(response)
|
||||
# Written beside the destination and moved in once the whole body has arrived, so
|
||||
# a transfer that fails part way leaves whatever was already there untouched and
|
||||
# never leaves a truncated file under the real name.
|
||||
fd, tmp_path = self._open_partial(path)
|
||||
try:
|
||||
with open(path, "wb") as f:
|
||||
with os.fdopen(fd, "wb") as f:
|
||||
for chunk in response.iter_content(chunk_size=4096):
|
||||
if chunk:
|
||||
f.write(chunk)
|
||||
f.flush()
|
||||
except PermissionError:
|
||||
print(f"Cannot create file `{path}`: Permission denied.")
|
||||
os.replace(tmp_path, path)
|
||||
except BaseException:
|
||||
# Only the scratch file goes: anything at the destination was not written here.
|
||||
if os.path.exists(tmp_path):
|
||||
os.unlink(tmp_path)
|
||||
raise
|
||||
return True
|
||||
|
||||
@staticmethod
|
||||
def _open_partial(path: str) -> tuple[int, str]:
|
||||
"""
|
||||
Create a scratch file beside `path` and return it open for writing.
|
||||
|
||||
Beside it, so moving the finished download into place is a rename within one
|
||||
directory. `0o666` rather than a private mode because the umask is what decided
|
||||
the permissions of a downloaded file before, and still should.
|
||||
"""
|
||||
base = f"{path}.{os.getpid()}-{threading.get_ident()}"
|
||||
for attempt in range(100):
|
||||
tmp_path = f"{base}-{attempt}.mdrspart"
|
||||
try:
|
||||
return os.open(tmp_path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o666), tmp_path
|
||||
except FileExistsError:
|
||||
continue
|
||||
except OSError as e:
|
||||
raise UnexpectedException(f"Cannot write `{path}`: {e}")
|
||||
raise UnexpectedException(f"Could not create a temporary file beside `{path}`.")
|
||||
|
||||
def _get_mime_type(self, path: str) -> str:
|
||||
mt = mimetypes.guess_type(path)
|
||||
if mt:
|
||||
|
||||
+120
-77
@@ -6,7 +6,7 @@ from unicodedata import normalize
|
||||
from pydantic.dataclasses import dataclass
|
||||
|
||||
from mdrsclient.api import FilesApi, FoldersApi
|
||||
from mdrsclient.exceptions import IllegalArgumentException, MDRSException, UnexpectedException
|
||||
from mdrsclient.exceptions import IllegalArgumentException, UnexpectedException
|
||||
from mdrsclient.models import File, Folder, Laboratory
|
||||
from mdrsclient.models.file import find_file
|
||||
from mdrsclient.settings import CONCURRENT
|
||||
@@ -27,7 +27,6 @@ class DownloadFileInfo:
|
||||
|
||||
@dataclass
|
||||
class DownloadContext:
|
||||
hasError: bool
|
||||
isSkipIfExists: bool
|
||||
files: list[DownloadFileInfo]
|
||||
|
||||
@@ -47,54 +46,67 @@ class Uploader:
|
||||
laboratory = self.client.find_laboratory(laboratory_name)
|
||||
folder = self.client.find_folder(laboratory, r_path)
|
||||
files = self.client.find_files(folder.id)
|
||||
infos: list[UploadFileInfo] = []
|
||||
if os.path.isdir(l_path):
|
||||
if not is_recursive:
|
||||
raise IllegalArgumentException(f"Cannot upload `{local_path}`: Is a directory.")
|
||||
folder_api = FoldersApi(self.client.connection)
|
||||
folder_map: dict[str, Folder] = {}
|
||||
folder_map[r_path] = folder
|
||||
files_map: dict[str, list[File]] = {}
|
||||
files_map[r_path] = files
|
||||
l_basename = os.path.basename(l_path)
|
||||
for dirpath, _, filenames in os.walk(l_path, followlinks=True):
|
||||
sub = l_basename if dirpath == l_path else os.path.join(l_basename, os.path.relpath(dirpath, l_path))
|
||||
d_dirname = os.path.join(r_path, sub)
|
||||
d_basename = os.path.basename(d_dirname)
|
||||
# prepare destination parent path
|
||||
d_parent_dirname = os.path.dirname(d_dirname)
|
||||
if folder_map.get(d_parent_dirname) is None:
|
||||
parent_folder = self.client.find_folder(laboratory, d_parent_dirname)
|
||||
folder_map[d_parent_dirname] = parent_folder
|
||||
parent_files = self.client.find_files(parent_folder.id)
|
||||
files_map[d_parent_dirname] = parent_files
|
||||
# prepare destination path
|
||||
if folder_map.get(d_dirname) is None:
|
||||
d_folder = folder_map[d_parent_dirname].find_sub_folder(d_basename)
|
||||
if d_folder is None:
|
||||
d_folder_id = folder_api.create(normalize("NFC", d_basename), folder_map[d_parent_dirname].id)
|
||||
else:
|
||||
d_folder_id = d_folder.id
|
||||
print(d_dirname)
|
||||
folder_map[d_dirname] = folder_api.retrieve(d_folder_id)
|
||||
files_map[d_dirname] = self.client.find_files(d_folder_id)
|
||||
if d_folder is None:
|
||||
folder_map[d_parent_dirname].sub_folders.append(folder_map[d_dirname])
|
||||
# register upload file list
|
||||
for filename in filenames:
|
||||
infos.append(
|
||||
UploadFileInfo(folder_map[d_dirname], files_map[d_dirname], os.path.join(dirpath, filename))
|
||||
)
|
||||
infos = self.__collect_directory_uploads(laboratory, r_path, l_path, folder, files)
|
||||
else:
|
||||
infos.append(UploadFileInfo(folder, files, l_path))
|
||||
self.__multiple_upload(infos, is_skip_if_exists)
|
||||
infos = [UploadFileInfo(folder, files, l_path)]
|
||||
if not self.__multiple_upload(infos, is_skip_if_exists):
|
||||
# One file failing is worth reporting on its own line, and worth the caller
|
||||
# hearing about: a batch that lost files is not a batch that succeeded.
|
||||
raise UnexpectedException("Some files failed to upload.")
|
||||
|
||||
def __multiple_upload(self, infos: list[UploadFileInfo], is_skip_if_exists: bool) -> None:
|
||||
def __collect_directory_uploads(
|
||||
self, laboratory: Laboratory, r_path: str, l_path: str, folder: Folder, files: list[File]
|
||||
) -> list[UploadFileInfo]:
|
||||
"""Mirror a local directory tree on the remote, and list the files to send into it."""
|
||||
infos: list[UploadFileInfo] = []
|
||||
folder_api = FoldersApi(self.client.connection)
|
||||
folder_map: dict[str, Folder] = {}
|
||||
folder_map[r_path] = folder
|
||||
files_map: dict[str, list[File]] = {}
|
||||
files_map[r_path] = files
|
||||
l_basename = os.path.basename(l_path)
|
||||
for dirpath, _, filenames in os.walk(l_path, followlinks=True):
|
||||
sub = l_basename if dirpath == l_path else os.path.join(l_basename, os.path.relpath(dirpath, l_path))
|
||||
d_dirname = os.path.join(r_path, sub)
|
||||
d_basename = os.path.basename(d_dirname)
|
||||
# prepare destination parent path
|
||||
d_parent_dirname = os.path.dirname(d_dirname)
|
||||
if folder_map.get(d_parent_dirname) is None:
|
||||
parent_folder = self.client.find_folder(laboratory, d_parent_dirname)
|
||||
folder_map[d_parent_dirname] = parent_folder
|
||||
parent_files = self.client.find_files(parent_folder.id)
|
||||
files_map[d_parent_dirname] = parent_files
|
||||
# prepare destination path
|
||||
if folder_map.get(d_dirname) is None:
|
||||
d_folder = folder_map[d_parent_dirname].find_sub_folder(d_basename)
|
||||
if d_folder is None:
|
||||
d_folder_id = folder_api.create(normalize("NFC", d_basename), folder_map[d_parent_dirname].id)
|
||||
else:
|
||||
d_folder_id = d_folder.id
|
||||
print(d_dirname)
|
||||
folder_map[d_dirname] = folder_api.retrieve(d_folder_id)
|
||||
files_map[d_dirname] = self.client.find_files(d_folder_id)
|
||||
if d_folder is None:
|
||||
folder_map[d_parent_dirname].sub_folders.append(folder_map[d_dirname])
|
||||
# register upload file list
|
||||
for filename in filenames:
|
||||
infos.append(
|
||||
UploadFileInfo(folder_map[d_dirname], files_map[d_dirname], os.path.join(dirpath, filename))
|
||||
)
|
||||
return infos
|
||||
|
||||
def __multiple_upload(self, infos: list[UploadFileInfo], is_skip_if_exists: bool) -> bool:
|
||||
"""Send every file, and report whether all of them arrived."""
|
||||
file_api = FilesApi(self.client.connection)
|
||||
with ThreadPoolExecutor(max_workers=CONCURRENT) as pool:
|
||||
pool.map(lambda x: self.__multiple_upload_worker(file_api, x, is_skip_if_exists), infos)
|
||||
results = pool.map(lambda x: self.__multiple_upload_worker(file_api, x, is_skip_if_exists), infos)
|
||||
# Consumed inside the block: the results are what carry each worker's verdict.
|
||||
return all(list(results))
|
||||
|
||||
def __multiple_upload_worker(self, file_api: FilesApi, info: UploadFileInfo, is_skip_if_exists: bool) -> None:
|
||||
def __multiple_upload_worker(self, file_api: FilesApi, info: UploadFileInfo, is_skip_if_exists: bool) -> bool:
|
||||
basename = os.path.basename(info.path)
|
||||
file = find_file(info.files, basename)
|
||||
try:
|
||||
@@ -103,8 +115,14 @@ class Uploader:
|
||||
elif not is_skip_if_exists or file.size != os.path.getsize(info.path):
|
||||
file_api.update(file, info.path)
|
||||
print(os.path.join(info.folder.path, basename))
|
||||
except MDRSException as e:
|
||||
print(f"Error: {e}")
|
||||
except Exception as e:
|
||||
# Everything, not just the exceptions the API layer raises: the batch verdict
|
||||
# is read now that the results are consumed, and a file vanishing between the
|
||||
# walk and the upload would otherwise end the whole run with a traceback and
|
||||
# throw away what every other file did.
|
||||
print(f"Failed: {info.path}: {e}")
|
||||
return False
|
||||
return True
|
||||
|
||||
|
||||
class Downloader:
|
||||
@@ -120,6 +138,21 @@ class Downloader:
|
||||
password: str | None = None,
|
||||
excludes: list[str] | None = None,
|
||||
) -> None:
|
||||
if not self.__download(remote_path, local_path, is_recursive, is_skip_if_exists, password, excludes):
|
||||
# Every failure has already been printed against the file it belongs to.
|
||||
# This is what makes the command as a whole end in failure.
|
||||
raise UnexpectedException("Some files failed to download.")
|
||||
|
||||
def __download(
|
||||
self,
|
||||
remote_path: str,
|
||||
local_path: str,
|
||||
is_recursive: bool,
|
||||
is_skip_if_exists: bool,
|
||||
password: str | None,
|
||||
excludes: list[str] | None,
|
||||
) -> bool:
|
||||
"""Fetch what the remote path names, and report whether every file arrived."""
|
||||
excludes_clean = excludes or []
|
||||
# Detect DOI path: "remote:10.xxxx/prefix.ID[/optional/sub/path]"
|
||||
path_component = remote_path.split(":", 1)[1] if ":" in remote_path else ""
|
||||
@@ -144,12 +177,11 @@ class Downloader:
|
||||
file = find_file(r_parent_files, r_basename)
|
||||
if file is not None:
|
||||
if self.__check_excludes(excludes_clean, laboratory, r_parent_folder, file):
|
||||
return
|
||||
context = DownloadContext(False, is_skip_if_exists, [])
|
||||
return True
|
||||
context = DownloadContext(is_skip_if_exists, [])
|
||||
l_path = os.path.join(l_dirname, r_basename)
|
||||
context.files.append(DownloadFileInfo(file, l_path))
|
||||
self.__multiple_download(context)
|
||||
return
|
||||
return self.__multiple_download(context)
|
||||
else:
|
||||
folder_simple = r_parent_folder.find_sub_folder(r_basename)
|
||||
if folder_simple is None:
|
||||
@@ -161,19 +193,17 @@ class Downloader:
|
||||
if not is_recursive:
|
||||
# Non-recursive: download only the files at the top level of the DOI folder.
|
||||
files = self.client.find_files(folder.id)
|
||||
context = DownloadContext(False, is_skip_if_exists, [])
|
||||
context = DownloadContext(is_skip_if_exists, [])
|
||||
for file in files:
|
||||
if self.__check_excludes(excludes_clean, laboratory, folder, file):
|
||||
continue
|
||||
l_path = os.path.join(l_dirname, file.name)
|
||||
context.files.append(DownloadFileInfo(file, l_path))
|
||||
self.__multiple_download(context)
|
||||
return
|
||||
return self.__multiple_download(context)
|
||||
folder_api = FoldersApi(self.client.connection)
|
||||
self.__multiple_download_pickup_recursive_files(
|
||||
return self.__multiple_download_pickup_recursive_files(
|
||||
folder_api, laboratory, folder.id, l_dirname, excludes_clean, is_skip_if_exists
|
||||
)
|
||||
return
|
||||
|
||||
remote, laboratory_name, r_path = self.client.parse_remote_host_with_path(remote_path)
|
||||
r_path = r_path.rstrip("/")
|
||||
@@ -189,11 +219,11 @@ class Downloader:
|
||||
file = find_file(r_parent_files, r_basename)
|
||||
if file is not None:
|
||||
if self.__check_excludes(excludes_clean, laboratory, r_parent_folder, file):
|
||||
return
|
||||
context = DownloadContext(False, is_skip_if_exists, [])
|
||||
return True
|
||||
context = DownloadContext(is_skip_if_exists, [])
|
||||
l_path = os.path.join(l_dirname, r_basename)
|
||||
context.files.append(DownloadFileInfo(file, l_path))
|
||||
self.__multiple_download(context)
|
||||
return self.__multiple_download(context)
|
||||
else:
|
||||
folder = r_parent_folder.find_sub_folder(r_basename)
|
||||
if folder is None:
|
||||
@@ -201,7 +231,7 @@ class Downloader:
|
||||
if not is_recursive:
|
||||
raise IllegalArgumentException(f"Cannot download `{r_path}`: Is a folder.")
|
||||
folder_api = FoldersApi(self.client.connection)
|
||||
self.__multiple_download_pickup_recursive_files(
|
||||
return self.__multiple_download_pickup_recursive_files(
|
||||
folder_api, laboratory, folder.id, l_dirname, excludes_clean, is_skip_if_exists
|
||||
)
|
||||
|
||||
@@ -213,47 +243,60 @@ class Downloader:
|
||||
basedir: str,
|
||||
excludes: list[str],
|
||||
is_skip_if_exists: bool,
|
||||
) -> None:
|
||||
context = DownloadContext(False, is_skip_if_exists, [])
|
||||
folder = folder_api.retrieve(folder_id)
|
||||
files = self.client.find_files(folder.id)
|
||||
) -> bool:
|
||||
context = DownloadContext(is_skip_if_exists, [])
|
||||
try:
|
||||
folder = folder_api.retrieve(folder_id)
|
||||
files = self.client.find_files(folder.id)
|
||||
except Exception as e:
|
||||
print(f"Failed: {basedir}: {e}")
|
||||
return False
|
||||
dirname = os.path.join(basedir, folder.name)
|
||||
if self.__check_excludes(excludes, laboratory, folder, None):
|
||||
return
|
||||
if not os.path.exists(dirname):
|
||||
os.makedirs(dirname)
|
||||
return True
|
||||
try:
|
||||
# `exist_ok` rather than a prior check: two workers can reach the same parent.
|
||||
os.makedirs(dirname, exist_ok=True)
|
||||
except OSError as e:
|
||||
# One folder the client cannot make locally is not a reason to abandon its
|
||||
# siblings, which is what this walk now promises.
|
||||
print(f"Failed: {dirname}: {e}")
|
||||
return False
|
||||
print(dirname)
|
||||
for file in files:
|
||||
if self.__check_excludes(excludes, laboratory, folder, file):
|
||||
continue
|
||||
path = os.path.join(dirname, file.name)
|
||||
context.files.append(DownloadFileInfo(file, path))
|
||||
self.__multiple_download(context)
|
||||
if context.hasError:
|
||||
raise UnexpectedException("Some files failed to download.")
|
||||
succeeded = self.__multiple_download(context)
|
||||
# A folder that lost a file is still a folder whose sub-folders the user asked
|
||||
# for, so the walk carries on and the verdict is collected for the caller.
|
||||
for sub_folder in folder.sub_folders:
|
||||
self.__multiple_download_pickup_recursive_files(
|
||||
if not self.__multiple_download_pickup_recursive_files(
|
||||
folder_api, laboratory, sub_folder.id, dirname, excludes, is_skip_if_exists
|
||||
)
|
||||
):
|
||||
succeeded = False
|
||||
return succeeded
|
||||
|
||||
def __multiple_download(self, context: DownloadContext) -> None:
|
||||
def __multiple_download(self, context: DownloadContext) -> bool:
|
||||
"""Fetch every file in the batch, and report whether all of them arrived."""
|
||||
file_api = FilesApi(self.client.connection)
|
||||
with ThreadPoolExecutor(max_workers=CONCURRENT) as pool:
|
||||
results = pool.map(
|
||||
lambda x: self.__multiple_download_worker(file_api, x, context.isSkipIfExists), context.files
|
||||
)
|
||||
hasError = next(filter(lambda x: x is False, results), None)
|
||||
if hasError is not None:
|
||||
context.hasError = True
|
||||
# Consumed inside the block, and in full: every worker's verdict counts, not
|
||||
# just the first refusal.
|
||||
return all(list(results))
|
||||
|
||||
def __multiple_download_worker(self, file_api: FilesApi, info: DownloadFileInfo, is_skip_if_exists: bool) -> bool:
|
||||
if not is_skip_if_exists or not os.path.exists(info.path) or info.file.size != os.path.getsize(info.path):
|
||||
try:
|
||||
file_api.download(info.file, info.path)
|
||||
except Exception:
|
||||
print(f"Failed: {info.path}")
|
||||
if os.path.isfile(info.path):
|
||||
os.remove(info.path)
|
||||
except Exception as e:
|
||||
# Nothing to clear up: a failed transfer writes only to its own scratch
|
||||
# file beside the destination, and removes that itself.
|
||||
print(f"Failed: {info.path}: {e}")
|
||||
return False
|
||||
print(info.path)
|
||||
return True
|
||||
|
||||
Reference in New Issue
Block a user