diff --git a/CHANGELOG.md b/CHANGELOG.md index 4140d14..87b89c5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,14 +6,18 @@ ### Features - **fsspec**: Improve multi-filesystem iRODS/HTTP integration and related CLI handling. +- **remote**: Add `remote upload` source URI support for `s3://` and `s3a://` (in addition to local paths) to upload directly from external S3 endpoints. +- **repository**: Enable upload path plumbing to accept non-local source filesystems in aggregate/file/lexis repository `put()` flows. ### Tests - **fsspec**: Expand unit tests for multi-filesystem behavior. +- **remote**: Add unit tests for unsupported source schemes and S3-backed source uploads in `remote_upload`. ### Documentation - **docs**: Update configuration docs and development process guidance for backlog and merge tracking. +- **remote**: Update `remote upload` CLI/docs examples and usage to include S3 URI sources. ## v5.1.0 (2026-01-21) diff --git a/docs/source/details/remote.md b/docs/source/details/remote.md index 318ecb0..9df0c7b 100644 --- a/docs/source/details/remote.md +++ b/docs/source/details/remote.md @@ -101,7 +101,7 @@ owi remote push "all/id=my-local-id" --datacenter it4i ## `remote upload` -Create a dataset from metadata and upload all files from a local directory, or upload into an existing dataset ID. +Create a dataset from metadata and upload all files from a local directory or S3 URI, or upload into an existing dataset ID. ### Usage @@ -125,7 +125,7 @@ owi remote upload DIRECTORY [OPTIONS] ### Modes and required flags - **Create + upload mode** - - Required: `DIRECTORY`, `--repository`, `--zone`, `--metadata-file` + - Required: `DIRECTORY` (local path, `file://`, `s3://`, `s3a://`), `--repository`, `--zone`, `--metadata-file` - Optional: `--collection`, `--access`, `--yes` - **Existing dataset upload mode** - Required: `DIRECTORY`, `--dataset-id` @@ -174,6 +174,12 @@ owi remote upload ./imprints \ --repository lexis ``` +**Upload from external S3 source:** +```bash +owi remote upload s3://user:pwd@endpoint/bucket/path \ + --dataset-id 450feb34-2bb4-11f0-9523-0242ac140003 +``` + **Update metadata and upload to existing dataset:** ```bash owi remote upload ./imprints \ diff --git a/owilix/cli/remote.py b/owilix/cli/remote.py index caa5f04..70b21a3 100644 --- a/owilix/cli/remote.py +++ b/owilix/cli/remote.py @@ -267,7 +267,10 @@ def push( @app.command() def upload( ctx: typer.Context, - directory: str = typer.Argument(..., help="Local directory whose files should be uploaded recursively"), + directory: str = typer.Argument( + ..., + help="Source directory path or URI (local path, file://, s3://, s3a://) to upload recursively", + ), repository: Optional[str] = typer.Option( None, "--repository", @@ -316,7 +319,7 @@ def upload( yes: bool = typer.Option(False, "--yes", "-y", help="Skip confirmation prompt"), ): """ - Upload files from a local directory to a remote dataset. + Upload files from a local directory or S3 URI to a remote dataset. Modes: 1. Create + upload: provide --metadata-file (no --dataset-id) @@ -329,6 +332,7 @@ def upload( owi remote upload ./data --dataset-id owi remote upload ./data --repository lexis --dataset-id owi remote upload ./data --dataset-id --update-metadata-from-file ./metadata.json + owi remote upload s3://user:pwd@endpoint/bucket/path --dataset-id """ cli_ctx: CLIContext = ctx.obj from owilix.core.tasks.remote import remote_upload diff --git a/owilix/core/repository/aggregate.py b/owilix/core/repository/aggregate.py index 773b8de..897cfcf 100644 --- a/owilix/core/repository/aggregate.py +++ b/owilix/core/repository/aggregate.py @@ -175,9 +175,7 @@ class AggregatedRepository: @measure_performance("upload", lambda self, dataset, local_path, local_file, filesystem=None: local_path) def put(self, dataset: Dataset, local_path: str, local_file: str, filesystem: fsspec.AbstractFileSystem = None) -> str: - if filesystem is not None: - raise NotImplementedError("Putting file to other filesystem not implemented yet.") - return dataset.repository.put(dataset, local_path, local_file) + return dataset.repository.put(dataset, local_path, local_file, filesystem=filesystem) @measure_performance("download", lambda self, dataset, file, local_path, filesystem=None: os.path.join(local_path, os.path.basename(file))) def get(self, dataset: Dataset, file: str, local_path: str, filesystem: fsspec.AbstractFileSystem = None): diff --git a/owilix/core/repository/file.py b/owilix/core/repository/file.py index a17bfb4..a05a26c 100644 --- a/owilix/core/repository/file.py +++ b/owilix/core/repository/file.py @@ -295,7 +295,10 @@ class FileBasedRepository(AbstractRepository): remote_dir = os.path.dirname(remote_file) if not self.fs.exists(remote_dir): self.fs.mkdirs(remote_dir) - self.fs.put(local_path, remote_file) + if filesystem is None: + self.fs.put(local_path, remote_file) + else: + self.copy_large_file(local_path, remote_file, filesystem, self.fs) return remote_file def copy_large_file(self, source_path: str, destination_path: str, source_fs, destination_fs): diff --git a/owilix/core/repository/lexis.py b/owilix/core/repository/lexis.py index 8a3dd97..f1204f2 100644 --- a/owilix/core/repository/lexis.py +++ b/owilix/core/repository/lexis.py @@ -405,10 +405,18 @@ class LexisRepository(AbstractRepository): # Upload file start_time = time.time() - file_size = os.path.getsize(local_path) + if filesystem is None: + file_size = os.path.getsize(local_path) + source_handle = open(local_path, "rb") + else: + try: + file_size = int(filesystem.size(local_path)) + except Exception: + file_size = 0 + source_handle = filesystem.open(local_path, "rb") chunk_size = 8 * 1024 * 1024 written = 0 - with open(local_path, "rb") as source: + with source_handle as source: first_chunk = source.read(chunk_size) if first_chunk == b"": result = self.fs.data_objects.write(b"", lpath=remote_path, truncate=1, append=0, offset=0) diff --git a/owilix/core/tasks/remote.py b/owilix/core/tasks/remote.py index acad6ee..beb5680 100644 --- a/owilix/core/tasks/remote.py +++ b/owilix/core/tasks/remote.py @@ -10,11 +10,13 @@ import re import time import subprocess import tempfile +import urllib.parse from collections import Counter from concurrent.futures import ThreadPoolExecutor, as_completed from typing import Optional, Dict, List, Any import duckdb +import fsspec from rich.console import Console from owilix.core.types import CommandResult @@ -36,6 +38,75 @@ def _get_tld_extractor(): return _TLD_EXTRACTOR +def _resolve_upload_source( + directory: str, +) -> tuple[str, Optional[fsspec.AbstractFileSystem], List[tuple[str, str]]]: + """ + Resolve a local directory path or supported URI into upload-ready file entries. + + Returns: + source_label: Human-readable source descriptor. + source_fs: Filesystem object for source files, or None for local filesystem. + source_files: List of tuples (source_path, relative_path). + """ + parsed = urllib.parse.urlparse(directory) + + # Local path mode (backward compatible). + if parsed.scheme in ("", "file"): + source_dir = os.path.abspath(os.path.expanduser(parsed.path if parsed.scheme == "file" else directory)) + if not os.path.isdir(source_dir): + raise ValueError(f"Directory not found: {source_dir}") + + source_files: List[tuple[str, str]] = [] + for root, _dirs, files in os.walk(source_dir): + for filename in files: + abs_path = os.path.join(root, filename) + rel_path = os.path.relpath(abs_path, source_dir).replace(os.sep, "/") + source_files.append((abs_path, rel_path)) + source_files.sort(key=lambda item: item[1]) + return source_dir, None, source_files + + if parsed.scheme not in ("s3", "s3a"): + raise ValueError( + f"Unsupported source scheme '{parsed.scheme}'. Supported: local path, file://, s3://, s3a://" + ) + + if not parsed.hostname: + raise ValueError("S3 source URI must include an endpoint host") + + endpoint = parsed.hostname + if parsed.port: + endpoint = f"{endpoint}:{parsed.port}" + endpoint_url = f"https://{endpoint}" + + source_root = parsed.path.lstrip("/") + if not source_root: + raise ValueError("S3 source URI must include bucket/path, e.g. s3://user:pwd@endpoint/bucket/path") + + key = urllib.parse.unquote(parsed.username) if parsed.username else None + secret = urllib.parse.unquote(parsed.password) if parsed.password else None + fs = fsspec.filesystem( + "s3", + key=key, + secret=secret, + anon=not bool(key and secret), + client_kwargs={"endpoint_url": endpoint_url}, + ) + + if not fs.isdir(source_root): + raise ValueError(f"Source directory not found on S3 endpoint: {source_root}") + + root_prefix = source_root.rstrip("/") + "/" + source_files = [] + for path in fs.find(source_root): + if fs.isdir(path): + continue + rel_path = path[len(root_prefix):] if path.startswith(root_prefix) else os.path.basename(path) + source_files.append((path, rel_path.replace("\\", "/"))) + source_files.sort(key=lambda item: item[1]) + return directory, fs, source_files + + def remote_pull( manager: Any, specifier: str, @@ -581,9 +652,10 @@ def remote_upload( if access_norm not in ("public", "project"): return CommandResult(success=False, msg="Access must be 'public' or 'project'") - source_dir = os.path.abspath(os.path.expanduser(directory)) - if not os.path.isdir(source_dir): - return CommandResult(success=False, msg=f"Directory not found: {source_dir}") + try: + source_label, source_fs, source_files = _resolve_upload_source(directory) + except Exception as e: + return CommandResult(success=False, msg=str(e)) metadata_payload = None if not dataset_id: @@ -708,16 +780,8 @@ def remote_upload( if not storage_name or not storage_resource: return CommandResult(success=False, msg=f"Selected storage is missing required fields: {selected_storage}") - local_files: List[tuple[str, str]] = [] - for root, _dirs, files in os.walk(source_dir): - for filename in files: - abs_path = os.path.join(root, filename) - rel_path = os.path.relpath(abs_path, source_dir).replace(os.sep, "/") - local_files.append((abs_path, rel_path)) - local_files.sort(key=lambda item: item[1]) - - if not local_files and metadata_update_payload is None: - return CommandResult(success=False, msg=f"No files found in directory: {source_dir}") + if not source_files and metadata_update_payload is None: + return CommandResult(success=False, msg=f"No files found in source: {source_label}") effective_access = access_norm if dataset_id and isinstance(existing_record, dict): @@ -732,6 +796,7 @@ def remote_upload( console.print(f" Access: {effective_access}") if collection_name: console.print(f" Collection: {collection_name}") + console.print(f" Source: {source_label}") if dataset_id: console.print(f" Dataset ID: {dataset_id} (existing)") if effective_access != access_norm: @@ -741,7 +806,7 @@ def remote_upload( console.print(f" Storage name: {storage_name}") console.print(f" Storage resource: {storage_resource}") console.print(f" Title: {title}") - console.print(f" Files: {len(local_files):,}") + console.print(f" Files: {len(source_files):,}") if metadata_update_payload is not None: console.print(" Metadata update: yes (from file)") @@ -854,16 +919,16 @@ def remote_upload( uploaded = 0 failed = 0 - if local_files: + if source_files: dataset_base_path = record.get("absolute_path") or dataset.path or "" with currentItemProgress() as progress: - upload_task = progress.add_task("Files uploaded", total=max(1, len(local_files)), current_item="Starting") - error_task = progress.add_task("Failed uploads", total=max(1, len(local_files)), current_item="None") - for abs_path, rel_path in local_files: + upload_task = progress.add_task("Files uploaded", total=max(1, len(source_files)), current_item="Starting") + error_task = progress.add_task("Failed uploads", total=max(1, len(source_files)), current_item="None") + for source_path, rel_path in source_files: try: target_path = f"{dataset_base_path.rstrip('/')}/{rel_path}" if dataset_base_path else rel_path progress.update(upload_task, current_item=f"Uploading {rel_path} -> {target_path}") - manager.remote_data.put(dataset, abs_path, rel_path) + manager.remote_data.put(dataset, source_path, rel_path, filesystem=source_fs) uploaded += 1 progress.update(upload_task, advance=1, current_item=f"Uploaded -> {target_path}") except Exception as e: @@ -877,7 +942,7 @@ def remote_upload( object={"datasetId": dataset_id, "uploaded": uploaded, "failed": failed}, msg=( f"{'Created dataset' if created_dataset else 'Using dataset'} {dataset_id}, " - f"but {failed}/{len(local_files)} uploads failed" + f"but {failed}/{len(source_files)} uploads failed" ), ) if created_dataset: diff --git a/tests/owilix/core/tasks/test_remote.py b/tests/owilix/core/tasks/test_remote.py index 72d94d2..cdfd502 100644 --- a/tests/owilix/core/tasks/test_remote.py +++ b/tests/owilix/core/tasks/test_remote.py @@ -916,6 +916,18 @@ class TestSummarizeLocalHostStats: # remote_upload # --------------------------------------------------------------------------- class TestRemoteUpload: + def test_rejects_unsupported_source_scheme(self): + manager = MagicMock() + result = remote_upload( + manager=manager, + directory="ftp://example.com/data", + dataset_id="ds-1", + yes=True, + console=MagicMock(), + ) + assert result.success is False + assert "Unsupported source scheme" in result.msg + def test_create_mode_requires_repository(self, tmp_path): data_dir = tmp_path / "data" data_dir.mkdir() @@ -1095,3 +1107,42 @@ class TestRemoteUpload: kwargs = repo.ddi_api.update_dataset_metadata.call_args.kwargs assert kwargs["dataset_id"] == "ds-2" assert "datacite" in kwargs["metadata"] + + @patch("owilix.core.tasks.remote.fsspec.filesystem") + def test_dataset_id_mode_uploads_from_s3_source(self, mock_filesystem): + source_fs = MagicMock() + + def _isdir(path): + return path == "bucket/path" + + source_fs.isdir.side_effect = _isdir + source_fs.find.return_value = ["bucket/path/folder/a.txt", "bucket/path/b.txt"] + mock_filesystem.return_value = source_fs + + repo = MagicMock() + repo.ddi_api = MagicMock() + repo.ddi_api.get_dataset_info_by_id.return_value = {"id": "ds-3", "absolute_path": "/irods/path/ds-3"} + + manager = MagicMock() + manager.remote_data.get_repo_names.return_value = ["lexis"] + manager.remote_data.get_single_repo.return_value = repo + + result = remote_upload( + manager=manager, + directory="s3://user:pwd@endpoint/bucket/path", + repository=None, + zone=None, + collection_name=None, + dataset_id="ds-3", + yes=True, + console=MagicMock(), + ) + + assert result.success is True + assert manager.remote_data.put.call_count == 2 + first_call = manager.remote_data.put.call_args_list[0] + second_call = manager.remote_data.put.call_args_list[1] + assert first_call.args[2] == "b.txt" + assert first_call.kwargs["filesystem"] is source_fs + assert second_call.args[2] == "folder/a.txt" + assert second_call.kwargs["filesystem"] is source_fs