Merge pull request #15090 from davelopez/fix_task_history_export_ftp

Fix export history to FTP using tasks
This commit is contained in:
Marius van den Beek
2022-12-06 14:51:49 +01:00
committed by GitHub
6 changed files with 125 additions and 26 deletions
+2 -1
View File
@@ -47,13 +47,14 @@ def stream_url_to_file(
file_sources: Optional["ConfiguredFileSources"] = None,
prefix: str = "gx_file_stream",
dir: Optional[str] = None,
user_context=None,
) -> str:
temp_name: str
if file_sources and file_sources.looks_like_uri(path):
file_source_path = file_sources.get_file_source_path(path)
with tempfile.NamedTemporaryFile(prefix=prefix, delete=False, dir=dir) as temp:
temp_name = temp.name
file_source_path.file_source.realize_to(file_source_path.path, temp_name)
file_source_path.file_source.realize_to(file_source_path.path, temp_name, user_context=user_context)
elif path.startswith("base64://"):
with tempfile.NamedTemporaryFile(prefix=prefix, delete=False, dir=dir) as temp:
temp_name = temp.name
+45 -13
View File
@@ -3,7 +3,9 @@ from typing import Optional
from galaxy import model
from galaxy.exceptions import RequestParameterInvalidException
from galaxy.jobs.manager import JobManager
from galaxy.managers.context import ProvidesUserContext
from galaxy.managers.histories import HistoryManager
from galaxy.managers.users import UserManager
from galaxy.model.scoped_session import galaxy_scoped_session
from galaxy.model.store import (
DirectoryModelExportStore,
@@ -32,6 +34,28 @@ from galaxy.web.short_term_storage import (
)
class ModelStoreUserContext(ProvidesUserContext):
def __init__(self, app: MinimalManagerApp, user: model.User) -> None:
self._app = app
self._user = user
@property
def app(self):
return self._app
@property
def url_builder(self):
raise NotImplementedError("URL builder not available in ModelStore context.")
def get_user(self):
return self._user
def set_user(self, user):
raise NotImplementedError("Cannot change user from ModelStore context.")
user = property(get_user, set_user)
class ModelStoreManager:
def __init__(
self,
@@ -40,12 +64,14 @@ class ModelStoreManager:
sa_session: galaxy_scoped_session,
job_manager: JobManager,
short_term_storage_monitor: ShortTermStorageMonitor,
user_manager: UserManager,
):
self._app = app
self._sa_session = sa_session
self._job_manager = job_manager
self._history_manager = history_manager
self._short_term_storage_monitor = short_term_storage_monitor
self._user_manager = user_manager
def setup_history_export_job(self, request: SetupHistoryExportJob):
history_id = request.history_id
@@ -116,11 +142,13 @@ class ModelStoreManager:
model_store_format = request.model_store_format
export_files = "symlink" if request.include_files else None
target_uri = request.target_uri
user_context = self._build_user_context(request.user.user_id)
with model.store.get_export_store_factory(
self._app,
model_store_format,
export_files=export_files,
bco_export_options=self._bco_export_options(request),
user_context=user_context,
)(target_uri) as export_store:
invocation = self._sa_session.query(model.WorkflowInvocation).get(request.invocation_id)
export_store.export_workflow_invocation(
@@ -142,9 +170,10 @@ class ModelStoreManager:
model_store_format = request.model_store_format
export_files = "symlink" if request.include_files else None
target_uri = request.target_uri
with model.store.get_export_store_factory(self._app, model_store_format, export_files=export_files)(
target_uri
) as export_store:
user_context = self._build_user_context(request.user.user_id)
with model.store.get_export_store_factory(
self._app, model_store_format, export_files=export_files, user_context=user_context
)(target_uri) as export_store:
if request.content_type == HistoryContentType.dataset:
hda = self._sa_session.query(model.HistoryDatasetAssociation).get(request.content_id)
export_store.add_dataset(hda)
@@ -158,9 +187,10 @@ class ModelStoreManager:
model_store_format = request.model_store_format
export_files = "symlink" if request.include_files else None
target_uri = request.target_uri
with model.store.get_export_store_factory(self._app, model_store_format, export_files=export_files)(
target_uri
) as export_store:
user_context = self._build_user_context(request.user.user_id)
with model.store.get_export_store_factory(
self._app, model_store_format, export_files=export_files, user_context=user_context
)(target_uri) as export_store:
history = self._history_manager.by_id(request.history_id)
export_store.export_history(
history, include_hidden=request.include_hidden, include_deleted=request.include_deleted
@@ -175,17 +205,13 @@ class ModelStoreManager:
history = self._sa_session.query(model.History).get(history_id)
else:
history = None
user_id = request.user.user_id
if user_id:
galaxy_user = self._sa_session.query(model.User).get(user_id)
else:
galaxy_user = None
user_context = self._build_user_context(request.user.user_id)
model_import_store = source_to_import_store(
request.source_uri,
self._app,
galaxy_user,
import_options,
model_store_format=request.model_store_format,
user_context=user_context,
)
create_new_history = history is None and not request.for_library
if create_new_history:
@@ -201,6 +227,11 @@ class ModelStoreManager:
)
return object_tracker
def _build_user_context(self, user_id: int):
user = self._user_manager.by_id(user_id)
user_context = ModelStoreUserContext(self._app, user)
return user_context
def create_objects_from_store(
app: MinimalManagerApp,
@@ -213,12 +244,13 @@ def create_objects_from_store(
discarded_data=ImportDiscardedDataType.FORCE,
allow_library_creation=for_library,
)
user_context = ModelStoreUserContext(app, galaxy_user) if galaxy_user is not None else None
model_import_store = source_to_import_store(
payload.store_content_uri or payload.store_dict,
app=app,
galaxy_user=galaxy_user,
import_options=import_options,
model_store_format=payload.model_store_format,
user_context=user_context,
)
create_new_history = history is None and not for_library
if create_new_history:
+22 -8
View File
@@ -50,7 +50,10 @@ from galaxy.exceptions import (
ObjectNotFound,
RequestParameterInvalidException,
)
from galaxy.files import ConfiguredFileSources
from galaxy.files import (
ConfiguredFileSources,
ProvidesUserFileSourcesUserContext,
)
from galaxy.files.uris import stream_url_to_file
from galaxy.model.mapping import GalaxyModelMapping
from galaxy.model.metadata import MetadataCollection
@@ -1812,6 +1815,7 @@ class DirectoryModelExportStore(ModelExportStore):
export_files: Optional[str] = None,
strip_metadata_files: bool = True,
serialize_jobs: bool = True,
user_context=None,
) -> None:
"""
:param export_directory: path to export directory. Will be created if it does not exist.
@@ -1836,6 +1840,7 @@ class DirectoryModelExportStore(ModelExportStore):
sessionless = True
security = IdEncodingHelper(id_secret="randomdoesntmatter")
self.user_context = ProvidesUserFileSourcesUserContext(user_context)
self.file_sources = file_sources
self.serialize_jobs = serialize_jobs
self.sessionless = sessionless
@@ -2500,7 +2505,7 @@ class BcoModelExportStore(WorkflowInvocationOnlyExportStore):
file_source_path = self.file_sources.get_file_source_path(self.file_source_uri)
file_source = file_source_path.file_source
assert os.path.exists(self.out_file)
file_source.write_from(file_source_path.path, self.out_file)
file_source.write_from(file_source_path.path, self.out_file, user_context=self.user_context)
def _core_biocompute_object_and_object_id(self) -> Tuple[BioComputeObjectCore, str]:
assert self.app # need app.security to do anything...
@@ -2725,7 +2730,7 @@ class ROCrateArchiveModelExportStore(DirectoryModelExportStore, WriteCrates):
file_source_path = self.file_sources.get_file_source_path(self.file_source_uri)
file_source = file_source_path.file_source
assert os.path.exists(rval), rval
file_source.write_from(file_source_path.path, rval)
file_source.write_from(file_source_path.path, rval, user_context=self.user_context)
shutil.rmtree(self.temp_output_dir)
@@ -2756,7 +2761,7 @@ class TarModelExportStore(DirectoryModelExportStore):
file_source_path = self.file_sources.get_file_source_path(self.file_source_uri)
file_source = file_source_path.file_source
assert os.path.exists(self.out_file)
file_source.write_from(file_source_path.path, self.out_file)
file_source.write_from(file_source_path.path, self.out_file, user_context=self.user_context)
shutil.rmtree(self.temp_output_dir)
@@ -2799,12 +2804,16 @@ class BagArchiveModelExportStore(BagDirectoryModelExportStore):
file_source_path = self.file_sources.get_file_source_path(self.file_source_uri)
file_source = file_source_path.file_source
assert os.path.exists(rval)
file_source.write_from(file_source_path.path, rval)
file_source.write_from(file_source_path.path, rval, user_context=self.user_context)
shutil.rmtree(self.temp_output_dir)
def get_export_store_factory(
app, download_format: str, export_files=None, bco_export_options: Optional[BcoExportOptions] = None
app,
download_format: str,
export_files=None,
bco_export_options: Optional[BcoExportOptions] = None,
user_context=None,
) -> Callable[[StrPath], ModelExportStore]:
export_store_class: Union[
Type[TarModelExportStore],
@@ -2816,6 +2825,7 @@ def get_export_store_factory(
"app": app,
"export_files": export_files,
"serialize_dataset_objects": False,
"user_context": user_context,
}
if download_format in ["tar.gz", "tgz"]:
export_store_class = TarModelExportStore
@@ -2871,10 +2881,11 @@ def imported_store_for_metadata(
def source_to_import_store(
source: Union[str, dict],
app: StoreAppProtocol,
galaxy_user: Optional[model.User],
import_options: Optional[ImportOptions],
model_store_format: Optional[ModelStoreFormat] = None,
user_context=None,
) -> ModelImportStore:
galaxy_user = user_context.user if user_context else None
if isinstance(source, dict):
if model_store_format is not None:
raise Exception(
@@ -2893,7 +2904,10 @@ def source_to_import_store(
if source_uri.startswith("file://"):
source_uri = source_uri[len("file://") :]
if "://" in source_uri:
source_uri = stream_url_to_file(source_uri, app.file_sources, prefix="gx_import_model_store")
user_context = ProvidesUserFileSourcesUserContext(user_context)
source_uri = stream_url_to_file(
source_uri, app.file_sources, prefix="gx_import_model_store", user_context=user_context
)
delete = True
target_path = source_uri
if target_path.endswith(".json"):
+21 -4
View File
@@ -37,9 +37,15 @@ def get_posix_file_source_config(root_dir: str, roles: str, groups: str, include
return rval
def create_file_source_config_file_on(temp_dir, root_dir, include_test_data_dir):
def create_file_source_config_file_on(
temp_dir,
root_dir,
include_test_data_dir,
required_role_expression,
required_group_expression,
):
file_contents = get_posix_file_source_config(
root_dir, REQUIRED_ROLE_EXPRESSION, REQUIRED_GROUP_EXPRESSION, include_test_data_dir
root_dir, required_role_expression, required_group_expression, include_test_data_dir
)
file_path = os.path.join(temp_dir, "file_sources_conf_posix.yml")
with open(file_path, "w") as f:
@@ -53,14 +59,25 @@ class PosixFileSourceSetup:
include_test_data_dir: ClassVar[bool] = False
@classmethod
def handle_galaxy_config_kwds(cls, config, clazz_=None):
def handle_galaxy_config_kwds(
cls,
config,
clazz_=None,
# Require role for access but do not require groups by default on every test to simplify them
required_role_expression=REQUIRED_ROLE_EXPRESSION,
required_group_expression="",
):
temp_dir = os.path.realpath(mkdtemp())
clazz_ = clazz_ or cls
clazz_._test_driver.temp_directories.append(temp_dir)
clazz_.root_dir = os.path.join(temp_dir, "root")
file_sources_config_file = create_file_source_config_file_on(
temp_dir, clazz_.root_dir, clazz_.include_test_data_dir
temp_dir,
clazz_.root_dir,
clazz_.include_test_data_dir,
required_role_expression,
required_group_expression,
)
config["file_sources_config_file"] = file_sources_config_file
@@ -1,3 +1,4 @@
import os
import tarfile
from uuid import uuid4
@@ -41,6 +42,14 @@ class TestImportExportHistoryViaTasksIntegration(
def handle_galaxy_config_kwds(cls, config):
PosixFileSourceSetup.handle_galaxy_config_kwds(config, cls)
UsesCeleryTasks.handle_galaxy_config_kwds(config)
cls.setup_ftp_config(config)
@classmethod
def setup_ftp_config(cls, config):
ftp_dir = cls.temp_config_dir("ftp")
os.makedirs(ftp_dir)
config["ftp_upload_dir"] = ftp_dir
config["ftp_upload_site"] = "ftp://cow.com"
def setUp(self):
super().setUp()
@@ -93,6 +102,23 @@ class TestImportExportHistoryViaTasksIntegration(
bai_metadata = imported_bam_details["meta_files"][0]
assert bai_metadata["file_type"] == "bam_index"
def test_import_export_ftp(self):
history_name = f"for_export_ftp_async_{uuid4()}"
history_id = self.dataset_populator.setup_history_for_export_testing(history_name)
model_store_format = "rocrate.zip"
target_uri = f"gxftp://history.{model_store_format}"
self.dataset_populator.export_history_to_uri_async(history_id, target_uri, model_store_format)
self.dataset_populator.import_history_from_uri_async(target_uri, model_store_format)
last_history = self._get("histories?limit=1").json()
assert len(last_history) == 1
imported_history = last_history[0]
imported_history_id = imported_history["id"]
assert imported_history_id != history_id
assert imported_history["name"] == history_name
class TestImportExportHistoryContentsViaTasksIntegration(IntegrationTestCase, UsesCeleryTasks):
dataset_populator: DatasetPopulator
@@ -18,6 +18,15 @@ from galaxy_test.driver.integration_setup import (
class TestPosixFileSourceIntegration(PosixFileSourceSetup, integration_util.IntegrationTestCase):
dataset_populator: DatasetPopulator
@classmethod
def handle_galaxy_config_kwds(cls, config):
PosixFileSourceSetup.handle_galaxy_config_kwds(
config,
cls,
required_role_expression=REQUIRED_ROLE_EXPRESSION,
required_group_expression=REQUIRED_GROUP_EXPRESSION,
)
def setUp(self):
super().setUp()
self._write_file_fixtures()