diff --git a/docs/features/implemented/remote-sync-sidecar-inventory-workflow.md b/docs/features/implemented/remote-sync-sidecar-inventory-workflow.md new file mode 100644 index 0000000..c01550f --- /dev/null +++ b/docs/features/implemented/remote-sync-sidecar-inventory-workflow.md @@ -0,0 +1,247 @@ +# Remote Sync Sidecar Inventory Workflow + +## Status + +Implemented in `feature/remote-sync-sidecar` for the v5.4.0 release. + +## Problem + +The original remote round-trip workflow was path-presence based: + +- `remote pull` downloaded files missing locally. +- `remote push` uploaded files missing remotely. +- A local edit to an existing file path was invisible to push unless the user forced a broad overwrite. + +That made the desired workflow unsafe: + +```bash +owi remote pull all/id= +# edit one local file +owi remote push all/id= +``` + +The remote service does not expose checksums or revision metadata, and no remote-side schema changes are available. The implementation therefore uses weak but fast file identity based on relative path, file size, and mtime. + +## User-Facing Workflow + +Pull or refresh a local baseline: + +```bash +owi remote pull all/id= --files "**/*" --yes +``` + +Check local status without remote reads: + +```bash +owi remote status all/id= --files "**/*" +``` + +Show changed file paths: + +```bash +owi remote status all/id= --files "**/*" --details +``` + +Refresh the remote comparison before classifying: + +```bash +owi remote status all/id= --files "**/*" --refresh-remote --details +``` + +Push local additions and modifications detected from the sidecar baseline: + +```bash +owi remote push all/id= +``` + +`remote push` does not propagate local deletions in this release. + +## Sidecar Files + +Each local dataset can contain two sync sidecars: + +- `.owi-sync.json`: small control file with source repository, baseline timestamp, baseline completeness, and comparison mode. +- `.owi-files.json.gz`: compressed file inventory baseline. + +The normal sibling metadata file remains untouched, so metadata reads do not need to load the file inventory. + +## Sidecar Schema + +`.owi-sync.json`: + +```json +{ + "schemaVersion": 1, + "source": { + "repository": "lexis", + "zone": "IT4ILexisV2", + "access": "public", + "datasetId": "abc-123", + "collectionName": "main" + }, + "baseline": { + "capturedAt": "2026-04-30T10:15:00Z", + "inventoryFile": ".owi-files.json.gz", + "complete": true + }, + "comparison": { + "mode": "size+mtime", + "mtimeUnit": "epoch_seconds" + } +} +``` + +`.owi-files.json.gz`: + +```json +{ + "schemaVersion": 1, + "capturedAt": "2026-04-30T10:15:00Z", + "datasetId": "abc-123", + "source": { + "repository": "lexis", + "zone": "IT4ILexisV2" + }, + "files": [ + { + "path": "README.md", + "size": 1024, + "mtime": 1776848873 + } + ], + "remoteFiles": [ + { + "path": "README.md", + "size": 1024, + "mtime": 1771317913 + } + ] +} +``` + +`files` is the local post-pull baseline. `remoteFiles` is the remote pull-time baseline. Keeping both avoids false local dirty states caused by download-time mtime differences. + +## Inventory API + +Repositories implement: + +```python +def files_inventory(dataset, files_glob=None) -> list[dict]: + ... +``` + +Entries are normalized to: + +```python +{"path": "relative/path", "size": 123, "mtime": 1776848873} +``` + +Implementations: + +- `FileBasedRepository.files_inventory()` reads local filesystem metadata only. +- `LexisRepository.files_inventory()` uses remote metadata listing through `fs.find(..., detail=True)` where available. +- `AggregatedRepository.files_inventory()` delegates to the owning repository. + +## Status States + +`remote status` reports: + +- `clean`: local inventory matches local baseline. +- `dirty-local`: local files were added or modified, or baseline files were removed. +- `remote-diverged`: current remote inventory differs from remote baseline. +- `conflict`: local and remote changed the same relative path. +- `unknown`: sidecar missing, unreadable, or baseline incomplete. +- `not-found`: the specifier contained `id` or `internalID`, but no local dataset matched. + +The compact count format is: + +```text ++ ~ - +``` + +For example: + +```text ++0 ~1 -0 +``` + +means no added files, one modified file, and no removed files. + +## Exclusions + +Always excluded: + +- `.owi-sync.json` +- `.owi-files.json.gz` +- Vim swap files such as `.README.md.swp` + +Excluded by default as generated files: + +- `stats.json` +- `changelog.json` + +`README.md` is tracked by default. + +Use `--include-generated` to include generated files in status checks: + +```bash +owi remote status all/id= --include-generated --details +``` + +## Pull Behavior + +After a successful pull, `remote pull` writes sidecars for each completed dataset. If any file download fails, the baseline for that dataset is not written. + +When a broad pull asks for all files and remote dataset metadata reports files but the remote file inventory returns zero files, the baseline is marked incomplete: + +```json +{ + "baseline": { + "complete": false, + "message": "remote metadata reports 2 files, but remote inventory returned 0 files" + } +} +``` + +`remote status` then reports `unknown` instead of falsely reporting `clean`. + +## Push Behavior + +When a complete sidecar baseline exists, `remote push` uses the local sidecar diff to select files to upload: + +- local added files are uploaded +- local modified files are uploaded +- local removed files are not propagated +- sync sidecars are never uploaded + +If no usable sidecar exists, push falls back to the previous path-presence behavior. + +After a successful sidecar-aware push, the local and remote baselines are refreshed. + +## Known Limitations + +- The comparison is not cryptographic; equal size and equal mtime are treated as equal content. +- Remote divergence detection is only as reliable as remote size and mtime metadata. +- Delete propagation is not implemented in this release. +- Conflict resolution is classification only; no automatic merge is attempted. + +## Tests + +Unit coverage includes: + +- timestamp normalization +- sidecar read/write +- inventory diffing and classification +- local, remote, and aggregate `files_inventory()` +- `remote pull` sidecar creation and incomplete-baseline detection +- `remote status` clean, dirty, conflict, not-found, details, and generated-file behavior +- `remote push` metadata update regressions, sidecar exclusion, sidecar-based modified-file upload, and baseline refresh +- `local free` dataset ID metadata fallback + +Live integration coverage: + +```bash +uv run pytest tests/owilix/core/repository/test_integration.py::TestIntegration::test_remote_sync_sidecar_status_workflow -q -m integration -s +``` + +The integration test pulls the small public dataset `29dd0066-ff87-11f0-ad38-02a47ca5d9fd` into an isolated temporary target, verifies clean status, tests local add/modify/delete detection, restores all local files, and verifies clean status again. diff --git a/docs/source/details/remote.md b/docs/source/details/remote.md index b1a3ec7..d3b4ee4 100644 --- a/docs/source/details/remote.md +++ b/docs/source/details/remote.md @@ -157,6 +157,74 @@ owi remote pull lexis:latest --files "**/*.parquet" owi remote pull lexis:latest --threads 8 ``` +### Sync Sidecars + +After a successful local pull, `owi remote pull` writes `.owi-sync.json` and `.owi-files.json.gz` into the dataset directory. These sidecars capture the pulled local file baseline and the remote file metadata baseline used by `remote status` and sidecar-aware `remote push`. + +If a broad pull sees an incomplete remote file inventory, for example remote metadata says the dataset has files but the listing returns none, the baseline is marked incomplete and `remote status` reports `unknown` instead of falsely reporting `clean`. + +--- + +## `remote status` + +Inspect the sidecar-based sync state of local datasets that were pulled from a remote repository. + +`remote status` compares the current local file inventory against the baseline written by `remote pull`. It can also refresh the remote file inventory and classify likely remote divergence without reading file contents. + +### Usage + +```bash +owi remote status [SPECIFIER] [OPTIONS] +``` + +### Options + +- `--files GLOB`: Limit status checks to matching files. +- `--refresh-remote`: Fetch current remote file metadata and compare it to the remote baseline captured during pull. +- `--details`: Print changed file paths grouped by added, modified, and removed. +- `--include-generated`: Include generated helper files such as `stats.json` and `changelog.json`. + +### States + +- `clean`: Local files match the local baseline. +- `dirty-local`: Local files were added, modified, or removed. +- `remote-diverged`: Current remote metadata differs from the remote baseline. +- `conflict`: The same relative path changed both locally and remotely. +- `unknown`: The sidecar baseline is missing, unreadable, or incomplete. +- `not-found`: The specifier contained an ID but no local dataset matched it. + +The compact count format is `+added ~modified -removed`. + +### Examples + +**Refresh a pulled dataset baseline:** +```bash +owi remote pull all/id= --files "**/*" --yes +``` + +**Check local dirty state:** +```bash +owi remote status all/id= --files "**/*" +``` + +**Show changed paths:** +```bash +owi remote status all/id= --files "**/*" --details +``` + +**Also compare against the current remote metadata:** +```bash +owi remote status all/id= --files "**/*" --refresh-remote --details +``` + +### Notes + +- The sync sidecars are `.owi-sync.json` and `.owi-files.json.gz` inside the local dataset directory. +- The comparison uses relative path, size, and mtime; it does not use checksums. +- `README.md` is tracked by default. +- `.owi-*` sidecars and editor swap files are always ignored. +- `stats.json` and `changelog.json` are ignored unless `--include-generated` is used. + --- ## `remote push` @@ -171,10 +239,14 @@ owi remote push [SPECIFIER] [OPTIONS] ### Options +- `--files GLOB`: Limit uploaded candidates to matching files. - `--datacenter DC`: Target datacenter (if not specified in specifier). - `--overwrite`: Overwrite remote files if they exist. +- `--mdupdate / --no-mdupdate`: Update remote metadata before uploading files. - `--yes` (`-y`): Skip confirmation. +When a complete sidecar baseline exists, `remote push` uploads the files classified as locally added or modified by `remote status`. Removed files are not propagated in this release. Without a usable sidecar baseline, push falls back to the previous path-presence behavior. + ### Examples **Push a local dataset to IT4I:** diff --git a/owilix/cli/remote.py b/owilix/cli/remote.py index e167311..8862308 100644 --- a/owilix/cli/remote.py +++ b/owilix/cli/remote.py @@ -516,6 +516,39 @@ def upload( raise typer.Exit(code=1) +@app.command() +def status( + ctx: typer.Context, + specifier: str = typer.Argument(..., help="Local dataset specifier to inspect"), + files: str = typer.Option("**/*", "--files", "-f", help="File glob pattern"), + refresh_remote: bool = typer.Option(False, "--refresh-remote", help="Refresh remote inventory before classifying status"), + details: bool = typer.Option(False, "--details", help="Show changed file paths"), + include_generated: bool = typer.Option(False, "--include-generated", help="Track generated files such as stats.json and changelog.json"), +): + """Show sidecar-based sync status for local datasets.""" + cli_ctx: CLIContext = ctx.obj + from owilix.core.tasks.remote import remote_status + + result = remote_status( + manager=cli_ctx.owi, + specifier=specifier, + files=files, + refresh_remote=refresh_remote, + details=details, + include_generated=include_generated, + console=cli_ctx.console, + show_table=cli_ctx.output_format not in ("json", "jsonl"), + ) + + if not result.success: + cli_ctx.console.print(f"[red]Status failed: {result.msg}[/red]") + raise typer.Exit(code=1) + + if cli_ctx.output_format in ("json", "jsonl"): + with OutputWriter(cli_ctx) as writer: + writer.write_records(result.object or []) + + @app.command() def diff( ctx: typer.Context, diff --git a/owilix/compat/local.py b/owilix/compat/local.py index ed26a4b..942d5d4 100644 --- a/owilix/compat/local.py +++ b/owilix/compat/local.py @@ -126,7 +126,9 @@ def free(self, specifier): datasets = self.list_local_datasets_by_specifier(specifier) self.console.print(f"found {len(datasets)} datasets") for dataset in datasets: - if self.autoyes or ask_yes_no(self.console, f"Delete dataset {dataset.internalID}/{dataset.title} from {dataset.path}"): + dataset_id = dataset.metadata.get("internalID") or dataset.metadata.get("id") or "unknown" + title = dataset.metadata.get("title", "Untitled") or "Untitled" + if self.autoyes or ask_yes_no(self.console, f"Delete dataset {dataset_id}/{title} from {dataset.path}"): try: dataset.repository.fs.rm(dataset.path+".json") dataset.repository.fs.rm(dataset.path, recursive=True) diff --git a/owilix/core/repository/aggregate.py b/owilix/core/repository/aggregate.py index 897cfcf..5be42ac 100644 --- a/owilix/core/repository/aggregate.py +++ b/owilix/core/repository/aggregate.py @@ -3,6 +3,8 @@ Aggregated Repository Aggregates multiple repositories under one API with performance-based selection. """ +from __future__ import annotations + import json import logging @@ -155,6 +157,9 @@ class AggregatedRepository: def files(self, dataset: Dataset, files_glob: str | Sequence[str] = None) -> list: return dataset.repository.files(dataset, files_glob) + def files_inventory(self, dataset: Dataset, files_glob: str | Sequence[str] = None) -> list: + return dataset.repository.files_inventory(dataset, files_glob) + def files_details(self, dataset: Dataset, files_glob: str | Sequence[str] = None, count_rows=False) -> list: return dataset.repository.files_details(dataset, files_glob, count_rows=count_rows) diff --git a/owilix/core/repository/base.py b/owilix/core/repository/base.py index ada037c..96408d6 100644 --- a/owilix/core/repository/base.py +++ b/owilix/core/repository/base.py @@ -9,6 +9,8 @@ Classes: RepoPerformanceStats: Holds performance metrics for a repository AbstractRepository: Abstract base class for all repositories """ +from __future__ import annotations + import logging import os @@ -243,6 +245,10 @@ class AbstractRepository(ABC): def files(self, dataset: Dataset, files_glob: str | Sequence[str] = None) -> list: raise NotImplementedError() + @abstractmethod + def files_inventory(self, dataset: Dataset, files_glob: str | Sequence[str] = None) -> list: + raise NotImplementedError() + @abstractmethod def delete(self, dataset: Dataset) -> None: raise NotImplementedError() diff --git a/owilix/core/repository/file.py b/owilix/core/repository/file.py index a05a26c..f979ee2 100644 --- a/owilix/core/repository/file.py +++ b/owilix/core/repository/file.py @@ -8,6 +8,8 @@ Classes: LocalRepository: Local filesystem repository S3FileBasedRepository: S3-based repository """ +from __future__ import annotations + import asyncio import json @@ -27,7 +29,8 @@ from owilix.core.repository.base import ( measure_performance, ) from owilix.core.models.dataset import Dataset -from owilix.core.utils import check_path_for_uuid_filename, fill_file_details +from owilix.core.sync import normalize_inventory_entry, to_epoch_seconds +from owilix.core.utils import check_path_for_uuid_filename, fill_file_details, rel_can_path logger = logging.getLogger("owilix") @@ -198,10 +201,14 @@ class FileBasedRepository(AbstractRepository): files_glob: str | Sequence[str] | None = None, ) -> Sequence[str]: """Return a list of files under the dataset path matching glob patterns.""" - # Use the actual dataset.path rather than reconstructing it, - # since the path may differ from what _get_path() computes - # (e.g., if collectionName in metadata differs from actual directory) - # TODO: Check LexisRepository.files() for the same issue with remote datasets + return [entry["absolute_path"] for entry in self.files_inventory(dataset, files_glob)] + + def files_inventory( + self, + dataset: Dataset, + files_glob: str | Sequence[str] | None = None, + ) -> list[dict]: + """Return file inventory entries using filesystem metadata only.""" dataset_path = dataset.path if dataset.path else self._get_path(dataset) if files_glob is None: @@ -214,7 +221,7 @@ class FileBasedRepository(AbstractRepository): def _is_absolute_like(p: str) -> bool: return ("://" in p or os.path.isabs(p) or p.startswith(dataset_path)) - hits: set[str] = set() + hits: dict[str, dict] = {} for pat in patterns: full_pattern = pat if _is_absolute_like(pat) else os.path.join(dataset_path, pat) @@ -222,13 +229,21 @@ class FileBasedRepository(AbstractRepository): if isinstance(result, dict): for path, info in result.items(): if info.get("type") in (None, "file"): - hits.add(path) + hits[path] = info else: for path in result: if not self.fs.isdir(path): - hits.add(path) - - return sorted(hits) + hits[path] = self.fs.info(path) + + inventory = [] + for abs_path, info in sorted(hits.items()): + inventory.append(normalize_inventory_entry( + path=rel_can_path(abs_path, dataset_path), + size=info.get("size", 0), + mtime=to_epoch_seconds(info.get("mtime") or info.get("modified")), + absolute_path=abs_path, + )) + return inventory def _get_collection_paths(self, access: str) -> List[str]: formatted_path = self.path.format(access=access) diff --git a/owilix/core/repository/lexis.py b/owilix/core/repository/lexis.py index 19e37b1..354613a 100644 --- a/owilix/core/repository/lexis.py +++ b/owilix/core/repository/lexis.py @@ -7,6 +7,8 @@ Uses: - OWILexisDatasetAPI for DDI metadata operations - Http2IrodsFileSystem for file operations via iRODS HTTP API """ +from __future__ import annotations + import logging import os @@ -19,6 +21,7 @@ import fsspec from owilix.core.repository.base import AbstractRepository from owilix.core.models.dataset import Dataset +from owilix.core.sync import normalize_inventory_entry, to_epoch_seconds logger = logging.getLogger("owilix") @@ -590,47 +593,71 @@ class LexisRepository(AbstractRepository): OPTIMIZED: Uses fs.find() which fetches entire tree in single HTTP request, then applies fnmatch patterns locally. """ + return [entry["absolute_path"] for entry in self.files_inventory(dataset, files_glob)] + + def files_inventory(self, dataset: Dataset, files_glob: str | Sequence[str] = None) -> list[dict]: + """Return file inventory entries using remote metadata only.""" path = self._get_path(dataset) - + if files_glob is None: patterns = ["**/*"] elif isinstance(files_glob, str): patterns = [files_glob] else: patterns = list(files_glob) - - # Use optimized find to get all files in one request + try: - all_entries = self.fs.find(path, withdirs=False) + all_entries = self.fs.find(path, withdirs=False, detail=True) except Exception as e: logger.warning(f"Optimized find failed: {e}, falling back to glob") - # Fallback to original implementation - hits = set() + hits = {} for pat in patterns: full_pattern = os.path.join(path, pat) try: - result = self.fs.glob(full_pattern) - for p in result: - if not self.fs.isdir(p): - hits.add(p) + result = self.fs.glob(full_pattern, detail=True) + if isinstance(result, dict): + for abs_path, info in result.items(): + if info.get("type") in (None, "file"): + hits[abs_path] = info + else: + for abs_path in result: + if not self.fs.isdir(abs_path): + hits[abs_path] = self.fs.info(abs_path) except Exception as e2: logger.warning(f"Glob pattern {full_pattern} failed: {e2}") - return sorted(hits) - - # Apply patterns locally (fast, no HTTP) - hits = set() + return [ + normalize_inventory_entry( + path=abs_path[len(path):].lstrip("/") if abs_path.startswith(path) else abs_path, + size=info.get("size", 0), + mtime=to_epoch_seconds(info.get("mtime") or info.get("modified")), + absolute_path=abs_path, + ) + for abs_path, info in sorted(hits.items()) + ] + + hits = {} for entry in all_entries: - # Get relative path for pattern matching - if entry.startswith(path): - rel_path = entry[len(path):].lstrip("/") + abs_path = entry.get("name", "") + if abs_path.startswith(path): + rel_path = abs_path[len(path):].lstrip("/") else: - rel_path = entry - + rel_path = abs_path + for pat in patterns: if _match_glob_pattern(rel_path, pat): - hits.add(entry) - - return sorted(hits) + hits[abs_path] = entry + break + + return [ + normalize_inventory_entry( + path=abs_path[len(path):].lstrip("/") if abs_path.startswith(path) else abs_path, + size=info.get("size", 0), + mtime=to_epoch_seconds(info.get("mtime") or info.get("modified")), + absolute_path=abs_path, + ctime=info.get("ctime"), + ) + for abs_path, info in sorted(hits.items()) + ] def files_details( self, diff --git a/owilix/core/sync.py b/owilix/core/sync.py new file mode 100644 index 0000000..8e42641 --- /dev/null +++ b/owilix/core/sync.py @@ -0,0 +1,168 @@ +"""Utilities for sidecar-based remote sync state.""" + +from __future__ import annotations + +import datetime as dt +import fnmatch +import gzip +import json +import os +from pathlib import Path +from typing import Any + + +def to_epoch_seconds(value: Any) -> int: + """Normalize various timestamp shapes to integer Unix epoch seconds.""" + if value is None or value == "": + return 0 + if isinstance(value, bool): + return int(value) + if isinstance(value, int): + return value + if isinstance(value, float): + return int(value) + if isinstance(value, dt.datetime): + if value.tzinfo is None: + return int(value.timestamp()) + return int(value.astimezone(dt.timezone.utc).timestamp()) + if isinstance(value, dt.date): + return int(dt.datetime(value.year, value.month, value.day, tzinfo=dt.timezone.utc).timestamp()) + + if isinstance(value, str): + stripped = value.strip() + if not stripped: + return 0 + try: + return int(float(stripped)) + except ValueError: + normalized = stripped.replace("Z", "+00:00") + try: + parsed = dt.datetime.fromisoformat(normalized) + except ValueError as exc: + raise ValueError(f"Unsupported mtime value: {value!r}") from exc + if parsed.tzinfo is None: + parsed = parsed.replace(tzinfo=dt.timezone.utc) + return int(parsed.timestamp()) + + raise TypeError(f"Unsupported mtime type: {type(value)!r}") + + +def normalize_inventory_entry( + path: str, + size: Any, + mtime: Any, + *, + absolute_path: str | None = None, + ctime: Any | None = None, +) -> dict[str, Any]: + """Create a normalized sync inventory record.""" + entry = { + "path": str(path).replace("\\", "/"), + "size": int(size or 0), + "mtime": to_epoch_seconds(mtime), + } + if absolute_path is not None: + entry["absolute_path"] = absolute_path + if ctime is not None: + entry["ctime"] = to_epoch_seconds(ctime) + return entry + + +def default_excludes(include_generated: bool = False) -> list[str]: + excludes = [ + ".owi-sync.json", + ".owi-files.json.gz", + ".*.swp", + ".*.swo", + "*.swp", + "*.swo", + ] + if not include_generated: + excludes.extend([ + "stats.json", + "changelog.json", + ]) + return excludes + + +def filter_inventory(inventory: list[dict[str, Any]], exclude: list[str] | None = None) -> list[dict[str, Any]]: + patterns = list(exclude or []) + if not patterns: + return list(inventory) + filtered = [] + for entry in inventory: + relpath = str(entry.get("path", "")) + if any(fnmatch.fnmatch(relpath, pattern) or fnmatch.fnmatch(os.path.basename(relpath), pattern) for pattern in patterns): + continue + filtered.append(entry) + return filtered + + +def diff_inventories(baseline: list[dict[str, Any]], current: list[dict[str, Any]]) -> dict[str, dict[str, Any]]: + baseline_map = {entry["path"]: {"size": int(entry.get("size", 0)), "mtime": int(entry.get("mtime", 0))} for entry in baseline} + current_map = {entry["path"]: {"size": int(entry.get("size", 0)), "mtime": int(entry.get("mtime", 0))} for entry in current} + + added = {path: current_map[path] for path in current_map.keys() - baseline_map.keys()} + removed = {path: baseline_map[path] for path in baseline_map.keys() - current_map.keys()} + modified = { + path: {"baseline": baseline_map[path], "current": current_map[path]} + for path in baseline_map.keys() & current_map.keys() + if baseline_map[path] != current_map[path] + } + unchanged = { + path: current_map[path] + for path in baseline_map.keys() & current_map.keys() + if baseline_map[path] == current_map[path] + } + + return { + "added": added, + "removed": removed, + "modified": modified, + "unchanged": unchanged, + } + + +def classify(local_diff: dict[str, dict[str, Any]], remote_diff: dict[str, dict[str, Any]] | None, baseline_complete: bool = True) -> str: + if not baseline_complete: + return "unknown" + + local_changed = bool(local_diff.get("added") or local_diff.get("removed") or local_diff.get("modified")) + if remote_diff is None: + return "dirty-local" if local_changed else "clean" + + remote_changed = bool(remote_diff.get("added") or remote_diff.get("removed") or remote_diff.get("modified")) + if local_changed and remote_changed: + local_paths = set(local_diff.get("added", {})) | set(local_diff.get("removed", {})) | set(local_diff.get("modified", {})) + remote_paths = set(remote_diff.get("added", {})) | set(remote_diff.get("removed", {})) | set(remote_diff.get("modified", {})) + if local_paths & remote_paths: + return "conflict" + if remote_changed: + return "remote-diverged" + if local_changed: + return "dirty-local" + return "clean" + + +def read_sync_sidecar(path: str | os.PathLike[str]) -> dict[str, Any]: + with open(path, "r", encoding="utf-8") as handle: + return json.load(handle) + + +def write_sync_sidecar(path: str | os.PathLike[str], data: dict[str, Any]) -> None: + target = Path(path) + target.parent.mkdir(parents=True, exist_ok=True) + with open(target, "w", encoding="utf-8") as handle: + json.dump(data, handle, indent=2, sort_keys=True) + + +def read_inventory_gz(path: str | os.PathLike[str]) -> dict[str, Any]: + with gzip.open(path, "rt", encoding="utf-8") as handle: + return json.load(handle) + + +def write_inventory_gz(path: str | os.PathLike[str], data: dict[str, Any]) -> None: + target = Path(path) + target.parent.mkdir(parents=True, exist_ok=True) + with gzip.open(target, "wt", encoding="utf-8") as handle: + json.dump(data, handle, sort_keys=True) diff --git a/owilix/core/tasks/local.py b/owilix/core/tasks/local.py index f7878c7..1dc0bcc 100644 --- a/owilix/core/tasks/local.py +++ b/owilix/core/tasks/local.py @@ -16,6 +16,15 @@ from owilix.core.utils import get_filesystem, fill_file_details, compare_filesys from owilix.cli._common.ui import input_dict, currentItemProgress, ask_yes_no +def _dataset_metadata_value(dataset: Any, key: str, default: Any = None) -> Any: + metadata = getattr(dataset, "metadata", None) + if metadata is not None and hasattr(metadata, "get"): + value = metadata.get(key, default) + if value is not None: + return value + return getattr(dataset, key, default) + + def remove_local_dataset( manager: Any, specifier: str, @@ -51,7 +60,13 @@ def remove_local_dataset( removed_count = 0 for dataset in datasets: - if auto_yes or ask_yes_no(console, f"Delete dataset {dataset.internalID}/{dataset.title} from {dataset.path}"): + dataset_id = ( + _dataset_metadata_value(dataset, "internalID") + or _dataset_metadata_value(dataset, "id") + or "unknown" + ) + title = _dataset_metadata_value(dataset, "title", "Untitled") or "Untitled" + if auto_yes or ask_yes_no(console, f"Delete dataset {dataset_id}/{title} from {dataset.path}"): try: # Remove json metadata file json_path = dataset.path + ".json" diff --git a/owilix/core/tasks/remote.py b/owilix/core/tasks/remote.py index 1cf5269..3ec3e62 100644 --- a/owilix/core/tasks/remote.py +++ b/owilix/core/tasks/remote.py @@ -4,8 +4,11 @@ Migrated from owilix/cmd/remote.py. """ import asyncio import inspect +import gzip import os import json +import fnmatch +from datetime import datetime, timezone import re import time import subprocess @@ -19,6 +22,16 @@ import duckdb import fsspec from rich.console import Console +from owilix.core.sync import ( + classify, + default_excludes, + diff_inventories, + filter_inventory, + read_inventory_gz, + read_sync_sidecar, + write_inventory_gz, + write_sync_sidecar, +) from owilix.core.types import CommandResult from owilix.core.db.models import OWIlixSQLQuery from owilix.core.db.duckdb_executor import OWIDuckDBSelectExecutor @@ -29,12 +42,96 @@ from owilix.cli._common.ui import currentItemProgress, ask_yes_no _TLD_EXTRACTOR = None +def _write_pull_sidecars( + local_dataset: Any, + remote_dataset: Any, + local_inventory: list[dict], + remote_inventory: list[dict] | None = None, + baseline_complete: bool = True, + baseline_message: str | None = None, +) -> None: + captured_at = datetime.now(timezone.utc).isoformat().replace("+00:00", "Z") + sync_path = os.path.join(local_dataset.path, ".owi-sync.json") + inventory_path = os.path.join(local_dataset.path, ".owi-files.json.gz") + + sync_payload = { + "schemaVersion": 1, + "source": { + "repository": (getattr(getattr(remote_dataset, "repository", None), "repo_name", None) if isinstance(getattr(getattr(remote_dataset, "repository", None), "repo_name", None), str) else None) or getattr(remote_dataset, "dataCenter", None), + "zone": getattr(remote_dataset, "zone", None) or remote_dataset.metadata.get("zone"), + "access": getattr(remote_dataset, "access", None) or remote_dataset.metadata.get("access"), + "datasetId": remote_dataset.metadata.get("id") or remote_dataset.metadata.get("internalID"), + "collectionName": remote_dataset.metadata.get("collectionName"), + }, + "baseline": { + "capturedAt": captured_at, + "inventoryFile": ".owi-files.json.gz", + "complete": baseline_complete, + }, + "comparison": {"mode": "size+mtime", "mtimeUnit": "epoch_seconds"}, + } + if baseline_message: + sync_payload["baseline"]["message"] = baseline_message + inventory_payload = { + "schemaVersion": 1, + "capturedAt": captured_at, + "datasetId": sync_payload["source"]["datasetId"], + "source": { + "repository": sync_payload["source"]["repository"], + "zone": sync_payload["source"]["zone"], + }, + "files": [ + {"path": entry["path"], "size": entry["size"], "mtime": entry["mtime"]} + for entry in local_inventory + ], + "remoteFiles": [ + {"path": entry["path"], "size": entry["size"], "mtime": entry["mtime"]} + for entry in (remote_inventory or local_inventory) + ], + } + + write_sync_sidecar(sync_path, sync_payload) + write_inventory_gz(inventory_path, inventory_payload) + + +def _is_broad_file_pattern(files_pattern: list[str]) -> bool: + return any(pattern in ("**/*", "*", "**") for pattern in files_pattern) + + +def _sync_path_matches(path: str, pattern: str) -> bool: + return fnmatch.fnmatch(path, pattern) or ( + pattern.startswith("**/") and fnmatch.fnmatch(path, pattern[3:]) + ) + + +def _sync_push_candidates(dataset: Any, local_inventory: list[dict], files_pattern: list[str]) -> set[str] | None: + sync_path, inventory_path = _sync_sidecar_paths(dataset) + if not os.path.exists(sync_path) or not os.path.exists(inventory_path): + return None + + try: + sync_payload = read_sync_sidecar(sync_path) + inventory_payload = read_inventory_gz(inventory_path) + except (OSError, json.JSONDecodeError, gzip.BadGzipFile): + return None + + if not sync_payload.get("baseline", {}).get("complete", False): + return None + + excludes = sync_payload.get("policy", {}).get("exclude") or default_excludes() + baseline_inventory = filter_inventory(list(inventory_payload.get("files", [])), excludes) + current_inventory = filter_inventory(local_inventory, excludes) + local_diff = diff_inventories(baseline_inventory, current_inventory) + candidates = set(local_diff.get("added", {})) | set(local_diff.get("modified", {})) + return {path for path in candidates if any(_sync_path_matches(path, pattern) for pattern in files_pattern)} + + def _get_tld_extractor(): global _TLD_EXTRACTOR if _TLD_EXTRACTOR is None: import tldextract # Keep default PSL behavior; instantiate once and reuse. - _TLD_EXTRACTOR = tldextract.TLDExtract(include_psl_private_domains=True) + _TLD_EXTRACTOR = tldextract.TLDExtract(include_psl_private_domains=True, cache_dir=os.path.join(tempfile.gettempdir(), "python-tldextract")) return _TLD_EXTRACTOR @@ -180,6 +277,7 @@ def remote_pull( files_pattern = parse_files_pattern(files) for d in datasets: + dataset_failed_files: List[tuple[str, str]] = [] console.print(f"Fetching files for {d.metadata.get('title', 'Unknown')} with file glob filter {files_pattern}") _remote_files = [rel_can_path(_p, d.path) for _p in manager.remote_data.files(d, files_pattern)] @@ -207,50 +305,71 @@ def remote_pull( if len(_missing_files) == 0: console.print(f"Dataset {d.metadata.get('title', 'Unknown')} is already up to date. Missing files are {len(_missing_files)}.") - continue - - with currentItemProgress() as _progress: - _dtask = _progress.add_task("Files downloaded", total=len(_missing_files), current_item="Setup") - _source_repo = d.repository + else: + with currentItemProgress() as _progress: + _dtask = _progress.add_task("Files downloaded", total=len(_missing_files), current_item="Setup") + _source_repo = d.repository - if num_threads > 1: - def download_one_file_threaded(rf): - try: - _source_repo.get( - d, - os.path.join(d.path, rf), - os.path.join(_ds_destination[0].path, os.path.dirname(rf)), - filesystem=None if push_to_remote is None else _dest_repo.fs - ) - _progress.update(_dtask, advance=1, current_item=f"Finished {rf}") - except Exception as e: - import traceback - traceback.print_exc() - console.print(f"[red]Error downloading file[/red] {rf}: {e.__class__.__name__}: {e}") - failed_files.append((rf, f"{e.__class__.__name__}: {e}")) - - with ThreadPoolExecutor(max_workers=num_threads) as executor: - future_to_file = {executor.submit(download_one_file_threaded, rf): rf for rf in _missing_files} - for future in as_completed(future_to_file): + if num_threads > 1: + def download_one_file_threaded(rf): try: - future.result() + _source_repo.get( + d, + os.path.join(d.path, rf), + os.path.join(_ds_destination[0].path, os.path.dirname(rf)), + filesystem=None if push_to_remote is None else _dest_repo.fs + ) + _progress.update(_dtask, advance=1, current_item=f"Finished {rf}") except Exception as e: - rf = future_to_file[future] - console.print(f"[red]Thread error downloading file[/red] {rf}: {e.__class__.__name__}: {e}") - failed_files.append((rf, f"{e.__class__.__name__}: {e}")) - else: - for _rf in _missing_files: - try: - _source_repo.get(d, os.path.join(d.path, _rf), - os.path.join(_ds_destination[0].path, os.path.dirname(_rf)), - filesystem=None if push_to_remote is None else _dest_repo.fs) - except Exception as e: - console.print(f"[red]Error downloading file[/red] {_rf}: {e.__class__.__name__}: {e}") - failed_files.append((_rf, f"{e.__class__.__name__}: {e}")) - _progress.update(_dtask, advance=1, current_item=f"Finished {_rf}") - - downloaded_count += len(_missing_files) - sum( - 1 for f, _ in failed_files if f in _missing_files + import traceback + traceback.print_exc() + console.print(f"[red]Error downloading file[/red] {rf}: {e.__class__.__name__}: {e}") + dataset_failed_files.append((rf, f"{e.__class__.__name__}: {e}")) + + with ThreadPoolExecutor(max_workers=num_threads) as executor: + future_to_file = {executor.submit(download_one_file_threaded, rf): rf for rf in _missing_files} + for future in as_completed(future_to_file): + try: + future.result() + except Exception as e: + rf = future_to_file[future] + console.print(f"[red]Thread error downloading file[/red] {rf}: {e.__class__.__name__}: {e}") + dataset_failed_files.append((rf, f"{e.__class__.__name__}: {e}")) + else: + for _rf in _missing_files: + try: + _source_repo.get(d, os.path.join(d.path, _rf), + os.path.join(_ds_destination[0].path, os.path.dirname(_rf)), + filesystem=None if push_to_remote is None else _dest_repo.fs) + except Exception as e: + console.print(f"[red]Error downloading file[/red] {_rf}: {e.__class__.__name__}: {e}") + dataset_failed_files.append((_rf, f"{e.__class__.__name__}: {e}")) + _progress.update(_dtask, advance=1, current_item=f"Finished {_rf}") + + downloaded_count += len(_missing_files) - len(dataset_failed_files) + + failed_files.extend(dataset_failed_files) + + if push_to_remote is None and not dataset_failed_files: + remote_inventory = d.repository.files_inventory(d, files_pattern) + local_inventory = _dest_repo.files_inventory(_ds_destination[0], files_pattern) + expected_file_count = int(d.metadata.get("fileCount", 0) or 0) + baseline_complete = True + baseline_message = None + if expected_file_count > 0 and not remote_inventory and _is_broad_file_pattern(files_pattern): + baseline_complete = False + baseline_message = ( + f"remote metadata reports {expected_file_count} files, " + "but remote inventory returned 0 files" + ) + console.print(f"[yellow]Warning:[/yellow] {baseline_message}") + _write_pull_sidecars( + _ds_destination[0], + d, + local_inventory, + remote_inventory, + baseline_complete=baseline_complete, + baseline_message=baseline_message, ) if failed_files: @@ -569,11 +688,18 @@ def remote_push( for d in datasets: title = d.metadata.get('title', 'Unknown') + dataset_id = d.metadata.get("internalID") or d.metadata.get("id") or getattr(d, "internalID", None) + collection_name = d.metadata.get("collectionName") console.print(f"Fetching files for {title}") - _local_files = [rel_can_path(_p, d.path) for _p in manager.local.files(d, files_pattern)] + local_inventory = manager.local.files_inventory(d, files_pattern) + _local_files = [ + entry["path"] + for entry in filter_inventory(local_inventory, default_excludes()) + ] + sync_candidates = None if overwrite else _sync_push_candidates(d, local_inventory, files_pattern) - _ds_remote = manager.remote_data.list(access=d.access, query={"internalID": d.internalID, - "collectionName": d.metadata.collectionName}) + _ds_remote = manager.remote_data.list(access=d.access, query={"internalID": dataset_id, + "collectionName": collection_name}) if len(_ds_remote) == 0: if manager.remote_data.exists_repo(dataCenter): @@ -586,13 +712,18 @@ def remote_push( console.print(f"Creating dataset {title} in datacenter {d.metadata.get('dataCenter')}") _ds_remote.append(manager.remote_data.create(d.metadata.get('dataCenter'), **d.metadata)) - manager.local.change_id(d, _ds_remote[-1].internalID) + manager.local.change_id(d, _dataset_identifier(_ds_remote[-1])) elif len(_ds_remote) > 1: - raise ValueError(f"Multiple datasets found for {d.internalID}. Dataset inconsistent") + raise ValueError(f"Multiple datasets found for {dataset_id}. Dataset inconsistent") else: if mdupdate: - _ds_remote[0].metadata.update({k: v for k, v in d.metadata.items() if k != "dataCenter"}) + metadata_update = { + k: v + for k, v in d.metadata.as_json_dict(version="V1").items() + if k != "dataCenter" + } + _ds_remote[0].metadata = metadata_update if hasattr(_ds_remote[0], 'reformat_metadata'): _ds_remote[0].reformat_metadata() manager.remote_data.update_metadata(_ds_remote[0]) @@ -608,7 +739,10 @@ def remote_push( _remote_files = [rel_can_path(_p, _ds_remote[0].path) for _p in manager.remote_data.files(_ds_remote[0], files_pattern)] - _missing_files = set(_local_files) - set(_remote_files) if not overwrite else set(_local_files) + if sync_candidates is not None: + _missing_files = sync_candidates + else: + _missing_files = set(_local_files) - set(_remote_files) if not overwrite else set(_local_files) if len(_missing_files) == 0: console.print(f"Dataset {d.metadata.get('title', 'Unknown')} is already up to date") @@ -617,8 +751,9 @@ def remote_push( if dataCenter is not None and _ds_remote[0].dataCenter != dataCenter: console.print(f"[yellow]Dataset {d.metadata.get('title', 'Unknown')} is already in datacenter {_ds_remote[0].dataCenter}. Requested {dataCenter} ignored.[/yellow]") - console.print(f"Uploading {len(_missing_files)} missing files now.") + console.print(f"Uploading {len(_missing_files)} changed files now.") + failed_uploads = 0 with currentItemProgress() as _progress: _dtask = _progress.add_task("Files uploaded", total=len(_missing_files), current_item="Setup") _error = _progress.add_task("Failed uploads", total=len(_missing_files), current_item="None") @@ -626,10 +761,16 @@ def remote_push( for _rf in _missing_files: _result = manager.remote_data.put(_ds_remote[0], os.path.join(d.path, _rf), _rf) if _result is None: + failed_uploads += 1 _progress.update(_error, advance=1, current_item=f"Failed {_rf}") else: _progress.update(_dtask, advance=1, current_item=f"Finished {_rf}") + if sync_candidates is not None and failed_uploads == 0: + refreshed_local_inventory = manager.local.files_inventory(d, files_pattern) + refreshed_remote_inventory = _ds_remote[0].repository.files_inventory(_ds_remote[0], files_pattern) + _write_pull_sidecars(d, _ds_remote[0], refreshed_local_inventory, refreshed_remote_inventory) + return CommandResult(success=True, object=datasets, msg=f"Pushed {len(datasets)} datasets") @@ -3957,6 +4098,224 @@ def summarize_local_host_stats( pass +def _sync_sidecar_paths(dataset: Any) -> tuple[str, str]: + return ( + os.path.join(dataset.path, ".owi-sync.json"), + os.path.join(dataset.path, ".owi-files.json.gz"), + ) + + +def _dataset_identifier(dataset: Any) -> str: + return dataset.metadata.get("internalID") or dataset.metadata.get("id") or getattr(dataset, "internalID", "") + + +def _dataset_title(dataset: Any) -> str: + return dataset.metadata.get("title", "Untitled") or "Untitled" + + +def _query_dataset_identifier(query: Any) -> str | None: + if not isinstance(query, dict): + return None + value = query.get("id") or query.get("internalID") + if isinstance(value, (list, tuple, set)): + return ",".join(str(item) for item in value) + return str(value) if value else None + + +def _status_record( + dataset_id: str, + title: str = "", + state: str = "unknown", + message: str = "", + *, + details: bool = False, +) -> dict[str, Any]: + record = { + "id": dataset_id, + "title": title, + "state": state, + "local_added": 0, + "local_modified": 0, + "local_removed": 0, + "remote_added": 0, + "remote_modified": 0, + "remote_removed": 0, + "message": message, + } + if details: + record["local_files"] = {"added": [], "modified": [], "removed": []} + record["remote_files"] = {"added": [], "modified": [], "removed": []} + return record + + +def _find_remote_dataset_for_status(manager: Any, source: dict[str, Any]) -> Any | None: + repo_name = source.get("repository") + repo = manager.remote_data.get_single_repo(repo_name) if repo_name else None + if repo is None and source.get("zone") and hasattr(manager.remote_data, "get_repo_for_zone"): + repo = manager.remote_data.get_repo_for_zone(source["zone"]) + if repo is None: + return None + + access = source.get("access") or "public" + dataset_id = source.get("datasetId") + collection_name = source.get("collectionName") + + queries = [] + if dataset_id and collection_name: + queries.append({"id": dataset_id, "collectionName": collection_name}) + queries.append({"internalID": dataset_id, "collectionName": collection_name}) + if dataset_id: + queries.append({"id": dataset_id}) + queries.append({"internalID": dataset_id}) + + for query in queries: + matches = repo.list(access=access, query=query) + if matches: + return matches[0] + return None + + +def remote_status( + manager: Any, + specifier: str, + files: str = "['**/*']", + refresh_remote: bool = False, + details: bool = False, + include_generated: bool = False, + console: Optional[Console] = None, + show_table: bool = True, +) -> CommandResult: + """Inspect sidecar-based sync state for local datasets.""" + if console is None: + console = Console() + + spec = manager.parse_specifier(specifier) + access, query = split_query_access(spec.get("query")) + local_datasets = manager.local.list( + access=access, + day=spec.get("day"), + duration=spec.get("duration") or 0, + query=query, + ) + files_pattern = parse_files_pattern(files) + + records = [] + if not local_datasets: + requested_id = _query_dataset_identifier(query) + if requested_id: + records.append( + _status_record( + requested_id, + state="not-found", + message="dataset not found locally", + details=details, + ) + ) + + for dataset in local_datasets: + dataset_id = _dataset_identifier(dataset) + sync_path, inventory_path = _sync_sidecar_paths(dataset) + record = _status_record(dataset_id, title=_dataset_title(dataset), details=details) + + if not os.path.exists(sync_path) or not os.path.exists(inventory_path): + record["message"] = "sync sidecar missing" + records.append(record) + continue + + try: + sync_payload = read_sync_sidecar(sync_path) + inventory_payload = read_inventory_gz(inventory_path) + except (OSError, json.JSONDecodeError, gzip.BadGzipFile) as exc: + record["message"] = f"failed to read sidecar: {exc}" + records.append(record) + continue + + excludes = sync_payload.get("policy", {}).get("exclude") or default_excludes( + include_generated=include_generated + ) + baseline_complete = bool(sync_payload.get("baseline", {}).get("complete", False)) + baseline_inventory = filter_inventory(list(inventory_payload.get("files", [])), excludes) + remote_baseline_inventory = filter_inventory( + list(inventory_payload.get("remoteFiles", inventory_payload.get("files", []))), + excludes, + ) + local_inventory = filter_inventory(manager.local.files_inventory(dataset, files_pattern), excludes) + local_diff = diff_inventories(baseline_inventory, local_inventory) + record["local_added"] = len(local_diff.get("added", {})) + record["local_modified"] = len(local_diff.get("modified", {})) + record["local_removed"] = len(local_diff.get("removed", {})) + if details: + record["local_files"] = { + "added": sorted(local_diff.get("added", {})), + "modified": sorted(local_diff.get("modified", {})), + "removed": sorted(local_diff.get("removed", {})), + } + + remote_diff = None + if refresh_remote: + remote_dataset = _find_remote_dataset_for_status(manager, sync_payload.get("source", {})) + if remote_dataset is None: + record["message"] = "remote dataset not found for baseline" + records.append(record) + continue + remote_inventory = filter_inventory(remote_dataset.repository.files_inventory(remote_dataset, files_pattern), excludes) + remote_diff = diff_inventories(remote_baseline_inventory, remote_inventory) + + record["state"] = classify(local_diff, remote_diff, baseline_complete=baseline_complete) + if remote_diff is not None: + record["remote_added"] = len(remote_diff.get("added", {})) + record["remote_modified"] = len(remote_diff.get("modified", {})) + record["remote_removed"] = len(remote_diff.get("removed", {})) + if details: + record["remote_files"] = { + "added": sorted(remote_diff.get("added", {})), + "modified": sorted(remote_diff.get("modified", {})), + "removed": sorted(remote_diff.get("removed", {})), + } + if not record["message"]: + record["message"] = sync_payload.get("baseline", {}).get("message") or ( + "baseline incomplete" if not baseline_complete else "ok" + ) + records.append(record) + + if show_table: + from rich.table import Table + + table = Table(title=f"Remote Sync Status: {specifier}") + table.add_column("state") + table.add_column("id") + table.add_column("title") + table.add_column("local") + if refresh_remote: + table.add_column("remote") + table.add_column("message") + + for record in records: + local_summary = f"+{record['local_added']} ~{record['local_modified']} -{record['local_removed']}" + row = [record["state"], record["id"], record["title"], local_summary] + if refresh_remote: + remote_summary = f"+{record['remote_added']} ~{record['remote_modified']} -{record['remote_removed']}" + row.append(remote_summary) + row.append(record["message"]) + table.add_row(*row) + console.print(table) + if details: + for record in records: + sections = [ + ("local", record.get("local_files", {})), + ("remote", record.get("remote_files", {})), + ] + for scope, changes in sections: + for kind in ("added", "modified", "removed"): + paths = changes.get(kind, []) + if paths: + console.print(f"\n[bold]{record['id']} {scope} {kind} ({len(paths)}):[/bold]") + for path in paths: + console.print(f" {path}") + + return CommandResult(success=True, object=records, msg=f"Status generated for {len(records)} datasets") + + def remote_diff( manager: Any, specifier: str, diff --git a/tests/owilix/cli/test_remote_status_cli.py b/tests/owilix/cli/test_remote_status_cli.py new file mode 100644 index 0000000..1ef060d --- /dev/null +++ b/tests/owilix/cli/test_remote_status_cli.py @@ -0,0 +1,54 @@ +from typer.testing import CliRunner +from unittest.mock import patch + +from owilix.cli import app +from owilix.core.types import CommandResult + + +runner = CliRunner() + + +def test_remote_status_json_output(tmp_path): + target = tmp_path / "owi" + target.mkdir() + + result_payload = CommandResult( + success=True, + object=[ + { + "id": "ds-001", + "title": "Dataset One", + "state": "clean", + "local_added": 0, + "local_modified": 0, + "local_removed": 0, + "remote_added": 0, + "remote_modified": 0, + "remote_removed": 0, + "message": "ok", + } + ], + msg="Status generated for 1 datasets", + ) + + with patch("owilix.core.tasks.remote.remote_status", return_value=result_payload) as remote_status: + result = runner.invoke( + app, + [ + "--target", + str(target), + "--format", + "json", + "remote", + "status", + "all", + "--details", + "--include-generated", + ], + ) + + assert result.exit_code == 0 + assert '"state": "clean"' in result.stdout + assert '"id": "ds-001"' in result.stdout + assert remote_status.call_args.kwargs["details"] is True + assert remote_status.call_args.kwargs["include_generated"] is True diff --git a/tests/owilix/core/repository/test_integration.py b/tests/owilix/core/repository/test_integration.py index 363eba3..c761476 100644 --- a/tests/owilix/core/repository/test_integration.py +++ b/tests/owilix/core/repository/test_integration.py @@ -12,6 +12,7 @@ The -s flag is important to see progress output during slow operations. import pytest import time import os +import json import yaml import shutil from typer.testing import CliRunner @@ -23,6 +24,55 @@ def elapsed(start: float) -> str: return f"[{time.time() - start:.1f}s]" +def _setup_isolated_target(target_dir: str) -> None: + """Create an isolated OWILIX config that reuses the developer credentials.""" + from owilix.core.manager import OWIlixManager + + mgr = OWIlixManager() + cfg = mgr.config.config.copy() + + if "repositories" not in cfg: + cfg["repositories"] = {} + if "config" not in cfg["repositories"]: + cfg["repositories"]["config"] = {} + if "local" not in cfg["repositories"]["config"]: + cfg["repositories"]["config"]["local"] = {"repository": "file", "options": {}} + + cfg["repositories"]["config"]["local"]["options"]["path"] = os.path.join(target_dir, "{access}") + + with open(os.path.join(target_dir, "owilix.cfg"), "w") as f: + yaml.dump(cfg, f) + + src_token_path = mgr._refresh_token_fn + if os.path.exists(src_token_path): + dst_token_dir = os.path.join(target_dir, ".tokens") + os.makedirs(dst_token_dir, exist_ok=True) + shutil.copy2(src_token_path, os.path.join(dst_token_dir, "refresh_token")) + + +def _json_status(runner: CliRunner, target_dir: str, specifier: str) -> dict: + result = runner.invoke( + app, + [ + "--target", + target_dir, + "--format", + "json", + "remote", + "status", + specifier, + "--files", + "**/*", + "--details", + ], + prog_name="owi", + ) + assert result.exit_code == 0, result.output + records = json.loads(result.output) + assert len(records) == 1 + return records[0] + + class TestIntegration: """Integration tests that verify CLI commands work end-to-end.""" @@ -152,6 +202,93 @@ class TestIntegration: print(f"{elapsed(start)} ✅ Pull test passed in {time.time() - start:.1f}s") + @pytest.mark.integration + def test_remote_sync_sidecar_status_workflow(self, runner, tmp_path, capsys): + """Exercise sidecar status on a small real dataset without touching remote data. + + The test pulls a small public dataset into an isolated local repository, verifies + clean sidecar status, then performs local-only add/modify/delete operations and + restores every local change before finishing. + """ + start = time.time() + target_dir = str(tmp_path) + dataset_id = "29dd0066-ff87-11f0-ad38-02a47ca5d9fd" + specifier = f"all/internalID={dataset_id}" + rel_file = "year=2026/month=1/day=31/language=ces/index.ciff.gz" + + print(f"\n{elapsed(start)} Starting sidecar workflow integration test...") + _setup_isolated_target(target_dir) + + result = runner.invoke( + app, + [ + "--target", + target_dir, + "--yes", + "remote", + "pull", + specifier, + "--files", + "**/*", + ], + prog_name="owi", + ) + assert result.exit_code == 0, result.output + + from owilix.core.manager import OWIlixManager + + mgr = OWIlixManager(owi_path=target_dir) + datasets = mgr.local.list(access="public", query={"internalID": dataset_id}) + assert len(datasets) == 1 + dataset_path = datasets[0].path + file_path = os.path.join(dataset_path, rel_file) + backup_path = os.path.join(target_dir, "index.ciff.gz.bak") + temp_path = os.path.join(dataset_path, "owilix-sync-test.txt") + + assert os.path.exists(os.path.join(dataset_path, ".owi-sync.json")) + assert os.path.exists(os.path.join(dataset_path, ".owi-files.json.gz")) + assert os.path.exists(file_path) + shutil.copy2(file_path, backup_path) + + try: + status = _json_status(runner, target_dir, specifier) + assert status["state"] == "clean" + assert (status["local_added"], status["local_modified"], status["local_removed"]) == (0, 0, 0) + + with open(temp_path, "w", encoding="utf-8") as f: + f.write("temporary sidecar test\n") + status = _json_status(runner, target_dir, specifier) + assert status["state"] == "dirty-local" + assert status["local_added"] == 1 + assert status["local_files"]["added"] == ["owilix-sync-test.txt"] + os.remove(temp_path) + + with open(file_path, "ab") as f: + f.write(b"x") + status = _json_status(runner, target_dir, specifier) + assert status["state"] == "dirty-local" + assert status["local_modified"] == 1 + assert status["local_files"]["modified"] == [rel_file] + shutil.copy2(backup_path, file_path) + + os.remove(file_path) + status = _json_status(runner, target_dir, specifier) + assert status["state"] == "dirty-local" + assert status["local_removed"] == 1 + assert status["local_files"]["removed"] == [rel_file] + finally: + if os.path.exists(backup_path): + os.makedirs(os.path.dirname(file_path), exist_ok=True) + shutil.copy2(backup_path, file_path) + os.remove(backup_path) + if os.path.exists(temp_path): + os.remove(temp_path) + + status = _json_status(runner, target_dir, specifier) + assert status["state"] == "clean" + assert (status["local_added"], status["local_modified"], status["local_removed"]) == (0, 0, 0) + print(f"{elapsed(start)} ✅ Sidecar workflow integration test passed") + @pytest.mark.integration def test_config_version(self, runner, capsys): """Test config version command. @@ -173,4 +310,3 @@ class TestIntegration: assert result.exit_code == 0 assert "Config version:" in result.output print(f"{elapsed(start)} ✅ Config version test passed") - diff --git a/tests/owilix/core/repository/test_repository_legacy_ported.py b/tests/owilix/core/repository/test_repository_legacy_ported.py index 4c68500..2d27ab9 100644 --- a/tests/owilix/core/repository/test_repository_legacy_ported.py +++ b/tests/owilix/core/repository/test_repository_legacy_ported.py @@ -168,6 +168,7 @@ class TinyRepo(AbstractRepository): # Implement abstract methods def exists(self, dataset): return True def files(self, dataset, files_glob=None): return [] + def files_inventory(self, dataset, files_glob=None): return [] def files_details(self, dataset, files_glob=None, count_rows=False): return [] def readlines(self, dataset, file_name): return "" def writelines(self, dataset, file_name, content): pass @@ -209,3 +210,38 @@ def test_aggregated_repo_picks_best_by_stats(tmp_path): loaded = agg.get_performance_stats() assert "fast" in loaded assert loaded["fast"]["list_bandwidth"]["count"] >= 1 + + +def test_filebased_files_inventory_returns_normalized_metadata(tmp_path): + repo = make_local_repo(tmp_path) + ds = make_dataset_for_repo(repo) + + nested = os.path.join(ds.path, "nested") + os.makedirs(nested, exist_ok=True) + file_path = os.path.join(nested, "data.txt") + with open(file_path, "wb") as fh: + fh.write(b"hello") + + inventory = repo.files_inventory(ds, "**/*.txt") + assert inventory == [ + { + "path": "nested/data.txt", + "size": 5, + "mtime": int(os.path.getmtime(file_path)), + "absolute_path": file_path, + } + ] + + +def test_aggregated_repo_delegates_files_inventory(tmp_path): + repo = make_local_repo(tmp_path) + ds = make_dataset_for_repo(repo) + file_path = os.path.join(ds.path, "data.bin") + with open(file_path, "wb") as fh: + fh.write(b"abc") + + agg = AggregatedRepository({"local": repo}) + inventory = agg.files_inventory(ds) + assert len(inventory) == 1 + assert inventory[0]["path"] == "data.bin" + assert inventory[0]["size"] == 3 diff --git a/tests/owilix/core/tasks/test_local.py b/tests/owilix/core/tasks/test_local.py new file mode 100644 index 0000000..6cba9f7 --- /dev/null +++ b/tests/owilix/core/tasks/test_local.py @@ -0,0 +1,56 @@ +from types import SimpleNamespace +from unittest.mock import MagicMock, patch + +import owilix.cli # noqa: F401 +from owilix.core.tasks.local import remove_local_dataset + + +def _manager_with_dataset(dataset): + manager = MagicMock() + manager.parse_specifier.return_value = { + "data_center": None, + "query": {"id": "ds-001"}, + "day": None, + "duration": 0, + } + manager.local.list.return_value = [dataset] + return manager + + +def _dataset_without_direct_metadata_attrs(path="/tmp/ds-001"): + fs = MagicMock() + fs.exists.return_value = True + repository = SimpleNamespace(fs=fs) + return SimpleNamespace( + metadata={ + "internalID": "ds-001", + "id": "ds-001", + "title": "Dataset One", + }, + path=path, + repository=repository, + ) + + +def test_remove_local_dataset_uses_metadata_id_with_auto_yes(): + dataset = _dataset_without_direct_metadata_attrs() + manager = _manager_with_dataset(dataset) + + result = remove_local_dataset(manager, "all/id=ds-001", console=MagicMock(), auto_yes=True) + + assert result.success is True + assert "Freed 1 datasets" in result.msg + dataset.repository.fs.rm.assert_any_call("/tmp/ds-001.json") + dataset.repository.fs.rm.assert_any_call("/tmp/ds-001", recursive=True) + + +@patch("owilix.core.tasks.local.ask_yes_no", return_value=True) +def test_remove_local_dataset_prompt_uses_metadata_values(mock_ask): + dataset = _dataset_without_direct_metadata_attrs() + manager = _manager_with_dataset(dataset) + + result = remove_local_dataset(manager, "all/id=ds-001", console=MagicMock(), auto_yes=False) + + assert result.success is True + prompt = mock_ask.call_args.args[1] + assert "ds-001/Dataset One" in prompt diff --git a/tests/owilix/core/tasks/test_remote.py b/tests/owilix/core/tasks/test_remote.py index 1abb30f..983a955 100644 --- a/tests/owilix/core/tasks/test_remote.py +++ b/tests/owilix/core/tasks/test_remote.py @@ -45,8 +45,10 @@ from owilix.core.tasks.remote import ( remote_pull, remote_push, remote_remove, + remote_status, ) from owilix.core.types import CommandResult +from owilix.core.sync import read_inventory_gz, read_sync_sidecar, write_inventory_gz, write_sync_sidecar # --------------------------------------------------------------------------- @@ -705,14 +707,19 @@ class TestRemotePull: @patch("owilix.core.tasks.remote.currentItemProgress") @patch("owilix.core.tasks.remote.ask_yes_no", return_value=True) - def test_pull_up_to_date(self, mock_ask, mock_progress): + def test_pull_up_to_date(self, mock_ask, mock_progress, tmp_path): ds = _make_mock_dataset() manager = _make_mock_manager(datasets=[ds]) # remote files = local files → nothing to download manager.remote_data.files.return_value = ["/data/ds1/file1.parquet"] - dest_ds = _make_mock_dataset(ds_id="ds-001", path="/local/ds1") + dest_path = tmp_path / "local" / "ds1" + dest_path.mkdir(parents=True) + dest_ds = _make_mock_dataset(ds_id="ds-001", path=str(dest_path)) manager.local.list.return_value = [dest_ds] - manager.local.files.return_value = ["/local/ds1/file1.parquet"] + manager.local.files.return_value = [str(dest_path / "file1.parquet")] + ds.repository.files_inventory.return_value = [ + {"path": "file1.parquet", "size": 1, "mtime": 1} + ] result = remote_pull( manager, specifier="dc1/public", files="['**/*']", console=MagicMock() @@ -785,16 +792,21 @@ class TestRemotePull: @patch("owilix.core.tasks.remote.currentItemProgress") @patch("owilix.core.tasks.remote.ask_yes_no", return_value=True) - def test_pull_all_succeed(self, mock_ask, mock_progress): + def test_pull_all_succeed(self, mock_ask, mock_progress, tmp_path): """H2: Successful pull returns success=True with no failed files.""" ds = _make_mock_dataset() manager = _make_mock_manager(datasets=[ds]) manager.remote_data.files.return_value = ["/data/ds1/file1.parquet"] - dest_ds = _make_mock_dataset(ds_id="ds-001", path="/local/ds1") + dest_path = tmp_path / "local" / "ds1" + dest_path.mkdir(parents=True) + dest_ds = _make_mock_dataset(ds_id="ds-001", path=str(dest_path)) manager.local.list.return_value = [dest_ds] manager.local.files.return_value = [] # file is missing → will download ds.repository.get.return_value = None # download succeeds + ds.repository.files_inventory.return_value = [ + {"path": "file1.parquet", "size": 1, "mtime": 1} + ] result = remote_pull( manager, specifier="dc1/public", @@ -805,6 +817,311 @@ class TestRemotePull: assert result.success is True assert "failed" not in result.msg.lower() + @patch("owilix.core.tasks.remote.currentItemProgress") + @patch("owilix.core.tasks.remote.ask_yes_no", return_value=True) + def test_pull_writes_sidecars_on_success(self, mock_ask, mock_progress, tmp_path): + ds = _make_mock_dataset() + manager = _make_mock_manager(datasets=[ds]) + manager.remote_data.files.return_value = ["/data/ds1/file1.parquet"] + dest_path = tmp_path / "local" / "ds1" + dest_path.mkdir(parents=True) + dest_ds = _make_mock_dataset(ds_id="ds-001", path=str(dest_path)) + manager.local.list.return_value = [dest_ds] + manager.local.files.return_value = [] + ds.repository.get.return_value = None + ds.repository.files_inventory.return_value = [ + {"path": "file1.parquet", "size": 7, "mtime": 12345, "absolute_path": "/data/ds1/file1.parquet"} + ] + manager.local.files_inventory.return_value = [ + {"path": "file1.parquet", "size": 7, "mtime": 99999, "absolute_path": str(dest_path / "file1.parquet")} + ] + + result = remote_pull(manager, specifier="dc1/public", files="['**/*']", auto_yes=True, console=MagicMock()) + + assert result.success is True + assert (dest_path / ".owi-sync.json").exists() + assert (dest_path / ".owi-files.json.gz").exists() + inventory = read_inventory_gz(dest_path / ".owi-files.json.gz") + assert inventory["files"][0]["mtime"] == 99999 + assert inventory["remoteFiles"][0]["mtime"] == 12345 + + @patch("owilix.core.tasks.remote.currentItemProgress") + @patch("owilix.core.tasks.remote.ask_yes_no", return_value=True) + def test_pull_marks_sidecar_incomplete_when_broad_remote_inventory_is_empty( + self, mock_ask, mock_progress, tmp_path + ): + ds = _make_mock_dataset(file_count=2) + manager = _make_mock_manager(datasets=[ds]) + manager.remote_data.files.return_value = [] + dest_path = tmp_path / "local" / "ds1" + dest_path.mkdir(parents=True) + dest_ds = _make_mock_dataset(ds_id="ds-001", path=str(dest_path)) + manager.local.list.return_value = [dest_ds] + manager.local.files.return_value = [] + ds.repository.files_inventory.return_value = [] + manager.local.files_inventory.return_value = [] + + result = remote_pull(manager, specifier="dc1/public", files="['**/*']", auto_yes=True, console=MagicMock()) + + assert result.success is True + sidecar = read_sync_sidecar(dest_path / ".owi-sync.json") + assert sidecar["baseline"]["complete"] is False + assert "remote metadata reports 2 files" in sidecar["baseline"]["message"] + + +# --------------------------------------------------------------------------- +# remote_status +# --------------------------------------------------------------------------- +class TestRemoteStatus: + def _write_sidecars(self, dataset_path, dataset_id="ds-001", complete=True): + write_sync_sidecar( + dataset_path / ".owi-sync.json", + { + "schemaVersion": 1, + "source": { + "repository": "dc1", + "zone": "zone1", + "access": "public", + "datasetId": dataset_id, + "collectionName": "main", + }, + "baseline": { + "capturedAt": "2026-04-21T10:15:00Z", + "inventoryFile": ".owi-files.json.gz", + "complete": complete, + }, + "comparison": {"mode": "size+mtime", "mtimeUnit": "epoch_seconds"}, + }, + ) + write_inventory_gz( + dataset_path / ".owi-files.json.gz", + { + "schemaVersion": 1, + "capturedAt": "2026-04-21T10:15:00Z", + "datasetId": dataset_id, + "source": {"repository": "dc1", "zone": "zone1"}, + "files": [{"path": "file1.parquet", "size": 7, "mtime": 100}], + "remoteFiles": [{"path": "file1.parquet", "size": 7, "mtime": 200}], + }, + ) + + def test_status_returns_not_found_row_for_missing_id(self): + manager = _make_mock_manager(local_datasets=[]) + manager.parse_specifier.return_value = { + "data_center": None, + "query": {"id": "missing-ds"}, + "day": None, + "duration": 0, + } + + result = remote_status(manager, specifier="all/id=missing-ds", console=MagicMock(), show_table=False) + + assert result.success is True + assert result.object == [ + { + "id": "missing-ds", + "title": "", + "state": "not-found", + "local_added": 0, + "local_modified": 0, + "local_removed": 0, + "remote_added": 0, + "remote_modified": 0, + "remote_removed": 0, + "message": "dataset not found locally", + } + ] + + def test_status_returns_not_found_row_for_missing_internal_id_with_details(self): + manager = _make_mock_manager(local_datasets=[]) + manager.parse_specifier.return_value = { + "data_center": None, + "query": {"internalID": "missing-internal"}, + "day": None, + "duration": 0, + } + + result = remote_status( + manager, + specifier="all/internalID=missing-internal", + details=True, + console=MagicMock(), + show_table=False, + ) + + assert result.success is True + assert result.object[0]["id"] == "missing-internal" + assert result.object[0]["state"] == "not-found" + assert result.object[0]["local_files"] == {"added": [], "modified": [], "removed": []} + assert result.object[0]["remote_files"] == {"added": [], "modified": [], "removed": []} + + def test_status_unknown_without_sidecars(self, tmp_path): + dataset_path = tmp_path / "local" / "ds1" + dataset_path.mkdir(parents=True) + ds = _make_mock_dataset(path=str(dataset_path)) + manager = _make_mock_manager(local_datasets=[ds]) + + result = remote_status(manager, specifier="all", console=MagicMock(), show_table=False) + + assert result.success is True + assert result.object[0]["state"] == "unknown" + assert result.object[0]["message"] == "sync sidecar missing" + + def test_status_clean(self, tmp_path): + dataset_path = tmp_path / "local" / "ds1" + dataset_path.mkdir(parents=True) + self._write_sidecars(dataset_path) + ds = _make_mock_dataset(path=str(dataset_path)) + manager = _make_mock_manager(local_datasets=[ds]) + manager.local.files_inventory.return_value = [{"path": "file1.parquet", "size": 7, "mtime": 100}] + + result = remote_status(manager, specifier="all", console=MagicMock(), show_table=False) + + assert result.success is True + assert result.object[0]["state"] == "clean" + assert result.object[0]["local_modified"] == 0 + + def test_status_reports_incomplete_baseline_message(self, tmp_path): + dataset_path = tmp_path / "local" / "ds1" + dataset_path.mkdir(parents=True) + self._write_sidecars(dataset_path, complete=False) + sync_path = dataset_path / ".owi-sync.json" + sidecar = read_sync_sidecar(sync_path) + sidecar["baseline"]["message"] = "remote metadata reports 2 files, but remote inventory returned 0 files" + write_sync_sidecar(sync_path, sidecar) + ds = _make_mock_dataset(path=str(dataset_path)) + manager = _make_mock_manager(local_datasets=[ds]) + manager.local.files_inventory.return_value = [{"path": "file1.parquet", "size": 7, "mtime": 100}] + + result = remote_status(manager, specifier="all", console=MagicMock(), show_table=False) + + assert result.success is True + assert result.object[0]["state"] == "unknown" + assert "remote metadata reports 2 files" in result.object[0]["message"] + + def test_status_dirty_local(self, tmp_path): + dataset_path = tmp_path / "local" / "ds1" + dataset_path.mkdir(parents=True) + self._write_sidecars(dataset_path) + ds = _make_mock_dataset(path=str(dataset_path)) + manager = _make_mock_manager(local_datasets=[ds]) + manager.local.files_inventory.return_value = [ + {"path": "file1.parquet", "size": 9, "mtime": 100}, + {"path": "file2.parquet", "size": 5, "mtime": 101}, + ] + + result = remote_status(manager, specifier="all", console=MagicMock(), show_table=False) + + assert result.success is True + assert result.object[0]["state"] == "dirty-local" + assert result.object[0]["local_added"] == 1 + assert result.object[0]["local_modified"] == 1 + + def test_status_tracks_readme_by_default_but_excludes_generated_files(self, tmp_path): + dataset_path = tmp_path / "local" / "ds1" + dataset_path.mkdir(parents=True) + self._write_sidecars(dataset_path) + ds = _make_mock_dataset(path=str(dataset_path)) + manager = _make_mock_manager(local_datasets=[ds]) + manager.local.files_inventory.return_value = [ + {"path": "file1.parquet", "size": 7, "mtime": 100}, + {"path": "README.md", "size": 10, "mtime": 100}, + {"path": "stats.json", "size": 20, "mtime": 100}, + ] + + result = remote_status(manager, specifier="all", details=True, console=MagicMock(), show_table=False) + + assert result.success is True + assert result.object[0]["state"] == "dirty-local" + assert result.object[0]["local_files"]["added"] == ["README.md"] + + def test_status_include_generated_tracks_generated_files(self, tmp_path): + dataset_path = tmp_path / "local" / "ds1" + dataset_path.mkdir(parents=True) + self._write_sidecars(dataset_path) + ds = _make_mock_dataset(path=str(dataset_path)) + manager = _make_mock_manager(local_datasets=[ds]) + manager.local.files_inventory.return_value = [ + {"path": "file1.parquet", "size": 7, "mtime": 100}, + {"path": "stats.json", "size": 20, "mtime": 100}, + ] + + result = remote_status( + manager, + specifier="all", + details=True, + include_generated=True, + console=MagicMock(), + show_table=False, + ) + + assert result.success is True + assert result.object[0]["state"] == "dirty-local" + assert result.object[0]["local_files"]["added"] == ["stats.json"] + + def test_status_remote_diverged(self, tmp_path): + dataset_path = tmp_path / "local" / "ds1" + dataset_path.mkdir(parents=True) + self._write_sidecars(dataset_path) + ds = _make_mock_dataset(path=str(dataset_path)) + manager = _make_mock_manager(local_datasets=[ds]) + manager.local.files_inventory.return_value = [{"path": "file1.parquet", "size": 7, "mtime": 100}] + + remote_repo = MagicMock() + remote_ds = _make_mock_dataset(path="/remote/ds1") + remote_ds.repository = remote_repo + remote_repo.list.return_value = [remote_ds] + remote_repo.files_inventory.return_value = [ + {"path": "file1.parquet", "size": 7, "mtime": 200}, + {"path": "file2.parquet", "size": 5, "mtime": 101}, + ] + manager.remote_data.get_single_repo.return_value = remote_repo + + result = remote_status(manager, specifier="all", refresh_remote=True, console=MagicMock(), show_table=False) + + assert result.success is True + assert result.object[0]["state"] == "remote-diverged" + assert result.object[0]["remote_added"] == 1 + + def test_status_details_include_changed_paths(self, tmp_path): + dataset_path = tmp_path / "local" / "ds1" + dataset_path.mkdir(parents=True) + self._write_sidecars(dataset_path) + ds = _make_mock_dataset(path=str(dataset_path)) + manager = _make_mock_manager(local_datasets=[ds]) + manager.local.files_inventory.return_value = [ + {"path": "file1.parquet", "size": 9, "mtime": 100}, + {"path": "file2.parquet", "size": 5, "mtime": 101}, + ] + + result = remote_status(manager, specifier="all", details=True, console=MagicMock(), show_table=False) + + assert result.success is True + assert result.object[0]["local_files"]["added"] == ["file2.parquet"] + assert result.object[0]["local_files"]["modified"] == ["file1.parquet"] + + def test_status_conflict(self, tmp_path): + dataset_path = tmp_path / "local" / "ds1" + dataset_path.mkdir(parents=True) + self._write_sidecars(dataset_path) + ds = _make_mock_dataset(path=str(dataset_path)) + manager = _make_mock_manager(local_datasets=[ds]) + manager.local.files_inventory.return_value = [{"path": "file1.parquet", "size": 8, "mtime": 100}] + + remote_repo = MagicMock() + remote_ds = _make_mock_dataset(path="/remote/ds1") + remote_ds.repository = remote_repo + remote_repo.list.return_value = [remote_ds] + remote_repo.files_inventory.return_value = [{"path": "file1.parquet", "size": 9, "mtime": 200}] + manager.remote_data.get_single_repo.return_value = remote_repo + + result = remote_status(manager, specifier="all", refresh_remote=True, console=MagicMock(), show_table=False) + + assert result.success is True + assert result.object[0]["state"] == "conflict" + assert result.object[0]["local_modified"] == 1 + assert result.object[0]["remote_modified"] == 1 + # --------------------------------------------------------------------------- # remote_push @@ -822,6 +1139,147 @@ class TestRemotePush: result = remote_push(manager, specifier="dc1/public", console=MagicMock()) assert result.success is True + def test_push_uses_metadata_id_without_direct_internal_id(self): + ds = _make_mock_dataset(path="/local/ds1") + delattr(ds, "internalID") + remote_ds = _make_mock_dataset(path="/remote/ds1") + manager = _make_mock_manager(local_datasets=[ds]) + manager.local.files_inventory.return_value = [{"path": "README.md", "size": 1, "mtime": 1}] + manager.remote_data.list.return_value = [remote_ds] + manager.remote_data.files.return_value = ["/remote/ds1/README.md"] + + result = remote_push(manager, specifier="dc1/public", mdupdate=False, auto_yes=True, console=MagicMock()) + + assert result.success is True + manager.remote_data.list.assert_called_with( + access="public", + query={"internalID": "ds-001", "collectionName": "main"}, + ) + + def test_push_does_not_upload_sync_sidecars(self): + ds = _make_mock_dataset(path="/local/ds1") + remote_ds = _make_mock_dataset(path="/remote/ds1") + manager = _make_mock_manager(local_datasets=[ds]) + manager.local.files_inventory.return_value = [ + {"path": ".owi-sync.json", "size": 1, "mtime": 1}, + {"path": ".owi-files.json.gz", "size": 1, "mtime": 1}, + {"path": ".README.md.swp", "size": 1, "mtime": 1}, + {"path": "README.md", "size": 1, "mtime": 1}, + ] + manager.remote_data.list.return_value = [remote_ds] + manager.remote_data.files.return_value = [] + manager.remote_data.put.return_value = "/remote/ds1/README.md" + + result = remote_push(manager, specifier="dc1/public", mdupdate=False, auto_yes=True, console=MagicMock()) + + assert result.success is True + manager.remote_data.put.assert_called_once_with(remote_ds, "/local/ds1/README.md", "README.md") + + def test_push_metadata_update_uses_json_dict_not_items(self): + ds = _make_mock_dataset(path="/local/ds1") + remote_ds = _make_mock_dataset(path="/remote/ds1") + manager = _make_mock_manager(local_datasets=[ds]) + manager.local.files_inventory.return_value = [{"path": "README.md", "size": 1, "mtime": 1}] + manager.remote_data.list.return_value = [remote_ds] + manager.remote_data.files.return_value = ["/remote/ds1/README.md"] + ds.metadata.items.side_effect = KeyError("endDate") + ds.metadata.as_json_dict.return_value = { + "title": "Dataset One", + "internalID": "ds-001", + "collectionName": "main", + "dataCenter": "local", + } + + result = remote_push(manager, specifier="dc1/public", mdupdate=True, auto_yes=True, console=MagicMock()) + + assert result.success is True + ds.metadata.as_json_dict.assert_called_once_with(version="V1") + manager.remote_data.update_metadata.assert_called_once_with(remote_ds) + + def test_push_metadata_update_replaces_remote_metadata_object(self): + ds = _make_mock_dataset(path="/local/ds1") + remote_ds = _make_mock_dataset(path="/remote/ds1") + manager = _make_mock_manager(local_datasets=[ds]) + manager.local.files_inventory.return_value = [{"path": "README.md", "size": 1, "mtime": 1}] + manager.remote_data.list.return_value = [remote_ds] + manager.remote_data.files.return_value = ["/remote/ds1/README.md"] + original_remote_metadata = remote_ds.metadata + ds.metadata.as_json_dict.return_value = { + "title": "Dataset One", + "internalID": "ds-001", + "collectionName": "main", + "last_modified_date": "2026-02-17 08:18:34", + "dataCenter": "local", + } + remote_ds.metadata.update.side_effect = AttributeError("'str' object has no attribute 'from_value'") + + result = remote_push(manager, specifier="dc1/public", mdupdate=True, auto_yes=True, console=MagicMock()) + + assert result.success is True + assert remote_ds.metadata.get("last_modified_date") == "2026-02-17 08:18:34" + assert remote_ds.metadata.get("dataCenter") is None + original_remote_metadata.update.assert_not_called() + manager.remote_data.update_metadata.assert_called_once_with(remote_ds) + + @patch("owilix.core.tasks.remote.currentItemProgress") + def test_push_uses_sidecar_modified_files_and_refreshes_baseline(self, mock_progress, tmp_path): + progress = MagicMock() + mock_progress.return_value.__enter__.return_value = progress + progress.add_task.return_value = 1 + + dataset_path = tmp_path / "local" / "ds1" + dataset_path.mkdir(parents=True) + readme = dataset_path / "README.md" + readme.write_text("changed\n") + ds = _make_mock_dataset(path=str(dataset_path)) + remote_ds = _make_mock_dataset(path="/remote/ds1") + manager = _make_mock_manager(local_datasets=[ds]) + manager.remote_data.list.return_value = [remote_ds] + manager.remote_data.files.return_value = ["/remote/ds1/README.md"] + manager.remote_data.put.return_value = "/remote/ds1/README.md" + manager.local.files_inventory.side_effect = [ + [{"path": "README.md", "size": 8, "mtime": 200, "absolute_path": str(readme)}], + [{"path": "README.md", "size": 8, "mtime": 200, "absolute_path": str(readme)}], + ] + remote_ds.repository.files_inventory.return_value = [ + {"path": "README.md", "size": 8, "mtime": 300, "absolute_path": "/remote/ds1/README.md"} + ] + write_sync_sidecar( + dataset_path / ".owi-sync.json", + { + "schemaVersion": 1, + "source": { + "repository": "dc1", + "zone": "zone1", + "access": "public", + "datasetId": "ds-001", + "collectionName": "main", + }, + "baseline": { + "capturedAt": "2026-04-21T10:15:00Z", + "inventoryFile": ".owi-files.json.gz", + "complete": True, + }, + "comparison": {"mode": "size+mtime", "mtimeUnit": "epoch_seconds"}, + }, + ) + write_inventory_gz( + dataset_path / ".owi-files.json.gz", + { + "schemaVersion": 1, + "files": [{"path": "README.md", "size": 7, "mtime": 100}], + "remoteFiles": [{"path": "README.md", "size": 7, "mtime": 100}], + }, + ) + + result = remote_push(manager, specifier="dc1/public", mdupdate=False, auto_yes=True, console=MagicMock()) + + assert result.success is True + manager.remote_data.put.assert_called_once_with(remote_ds, str(readme), "README.md") + refreshed = read_inventory_gz(dataset_path / ".owi-files.json.gz") + assert refreshed["files"] == [{"path": "README.md", "size": 8, "mtime": 200}] + assert refreshed["remoteFiles"] == [{"path": "README.md", "size": 8, "mtime": 300}] + # --------------------------------------------------------------------------- # remote_remove diff --git a/tests/owilix/core/test_repositories.py b/tests/owilix/core/test_repositories.py index 408d7a2..9d5e6df 100644 --- a/tests/owilix/core/test_repositories.py +++ b/tests/owilix/core/test_repositories.py @@ -205,6 +205,7 @@ class TinyRepo(AbstractRepository): def readlines(self, dataset, file_name: str): return "" def writelines(self, dataset, file_name: str, content: str): return None def files(self, dataset, files_glob: str | Sequence[str] = None): return [] + def files_inventory(self, dataset, files_glob: str | Sequence[str] = None): return [] def delete(self, dataset): return None def delete_by_id(self, id: str): return None def put(self, dataset, local_path: str, rel_filename: str, filesystem=None): return "" diff --git a/tests/owilix/core/test_sync.py b/tests/owilix/core/test_sync.py new file mode 100644 index 0000000..ba9e973 --- /dev/null +++ b/tests/owilix/core/test_sync.py @@ -0,0 +1,118 @@ +import datetime as dt + +from owilix.core.sync import ( + classify, + default_excludes, + diff_inventories, + filter_inventory, + normalize_inventory_entry, + read_inventory_gz, + read_sync_sidecar, + to_epoch_seconds, + write_inventory_gz, + write_sync_sidecar, +) + + +def test_to_epoch_seconds_accepts_common_shapes(): + assert to_epoch_seconds(1745143915) == 1745143915 + assert to_epoch_seconds(1745143915.9) == 1745143915 + assert to_epoch_seconds("1745143915") == 1745143915 + assert to_epoch_seconds("2026-04-21T10:15:15+00:00") == 1776766515 + assert to_epoch_seconds(dt.datetime(2026, 4, 21, 10, 15, 15, tzinfo=dt.timezone.utc)) == 1776766515 + + +def test_normalize_inventory_entry_normalizes_types(): + entry = normalize_inventory_entry( + path="a\\b.txt", + size="12", + mtime="1745143915", + absolute_path="/tmp/a/b.txt", + ctime="1745143900", + ) + assert entry == { + "path": "a/b.txt", + "size": 12, + "mtime": 1745143915, + "absolute_path": "/tmp/a/b.txt", + "ctime": 1745143900, + } + + +def test_sidecar_roundtrip(tmp_path): + sync_path = tmp_path / ".owi-sync.json" + inventory_path = tmp_path / ".owi-files.json.gz" + + sync_payload = {"schemaVersion": 1, "baseline": {"complete": True}} + inventory_payload = {"files": [{"path": "x", "size": 1, "mtime": 2}]} + + write_sync_sidecar(sync_path, sync_payload) + write_inventory_gz(inventory_path, inventory_payload) + + assert read_sync_sidecar(sync_path) == sync_payload + assert read_inventory_gz(inventory_path) == inventory_payload + + +def test_filter_inventory_excludes_internal_files(): + inventory = [ + {"path": ".owi-sync.json", "size": 1, "mtime": 1}, + {"path": ".README.md.swp", "size": 1, "mtime": 1}, + {"path": "README.md", "size": 1, "mtime": 1}, + {"path": "changelog.json", "size": 1, "mtime": 1}, + {"path": "data/file.parquet", "size": 2, "mtime": 2}, + ] + filtered = filter_inventory(inventory, default_excludes()) + assert filtered == [ + {"path": "README.md", "size": 1, "mtime": 1}, + {"path": "data/file.parquet", "size": 2, "mtime": 2}, + ] + + +def test_default_excludes_can_include_generated_files(): + inventory = [ + {"path": ".owi-sync.json", "size": 1, "mtime": 1}, + {"path": "stats.json", "size": 1, "mtime": 1}, + {"path": "changelog.json", "size": 1, "mtime": 1}, + {"path": "README.md", "size": 1, "mtime": 1}, + ] + filtered = filter_inventory(inventory, default_excludes(include_generated=True)) + assert filtered == [ + {"path": "stats.json", "size": 1, "mtime": 1}, + {"path": "changelog.json", "size": 1, "mtime": 1}, + {"path": "README.md", "size": 1, "mtime": 1}, + ] + + +def test_diff_inventories_detects_all_change_kinds(): + baseline = [ + {"path": "same.txt", "size": 1, "mtime": 10}, + {"path": "mtime.txt", "size": 1, "mtime": 10}, + {"path": "size.txt", "size": 1, "mtime": 10}, + {"path": "removed.txt", "size": 1, "mtime": 10}, + ] + current = [ + {"path": "same.txt", "size": 1, "mtime": 10}, + {"path": "mtime.txt", "size": 1, "mtime": 11}, + {"path": "size.txt", "size": 2, "mtime": 10}, + {"path": "added.txt", "size": 3, "mtime": 12}, + ] + + diff = diff_inventories(baseline, current) + assert set(diff["added"]) == {"added.txt"} + assert set(diff["removed"]) == {"removed.txt"} + assert set(diff["modified"]) == {"mtime.txt", "size.txt"} + assert set(diff["unchanged"]) == {"same.txt"} + + +def test_classify_states(): + empty = {"added": {}, "removed": {}, "modified": {}, "unchanged": {}} + changed = {"added": {"x": {}}, "removed": {}, "modified": {}, "unchanged": {}} + remote_other = {"added": {"y": {}}, "removed": {}, "modified": {}, "unchanged": {}} + remote_same = {"added": {"x": {}}, "removed": {}, "modified": {}, "unchanged": {}} + + assert classify(empty, None) == "clean" + assert classify(changed, None) == "dirty-local" + assert classify(empty, remote_other) == "remote-diverged" + assert classify(changed, remote_other) == "remote-diverged" + assert classify(changed, remote_same) == "conflict" + assert classify(empty, None, baseline_complete=False) == "unknown"