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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions .github/ISSUE_TEMPLATE/new_issue.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,15 +7,15 @@ assignees: ''

---

## Description
### Summary
A clear and concise description of what the problem is.

## Considerations
### Considerations
Make sure not to forget about these.

## Sub-tasks
### Sub-tasks
List any sub-tasks that might help create a clear flow or sub-issues.

## Definition of Done
### Definition of Done
This might help in cases when what is to be delivered is more than just one
piece of functional code.
38 changes: 38 additions & 0 deletions cfa/dataops/catalog.py
Original file line number Diff line number Diff line change
Expand Up @@ -300,6 +300,44 @@ def _get_version_blobs(self, version: str = "latest") -> list:
key=lambda x: x["creation_time"],
)

def download_version_to_local(
self, local_path: str, version: str = "latest", force: bool = False
) -> bool:
"""Download a specific version of the data to a local path

Args:
local_path (str): the local path to download to
version (str, optional): the version to download. Defaults to "latest".
force (bool, optional): whether to force re-download if local.
Returns:
bool: whether any files were written
"""
written = False
blobs = self._get_version_blobs(version)
for blob in blobs:
blob_data = read_blob_stream(
blob_url=blob["name"],
account_name=self.account,
container_name=self.container,
)
relative_path = blob["name"].removeprefix(f"{self.prefix}/")
local_file_path = os.path.join(local_path, relative_path)
local_dir = os.path.dirname(local_file_path)
os.makedirs(local_dir, exist_ok=True)
if os.path.exists(local_file_path) and not force:
continue
# Handle both raw bytes and objects with content_as_bytes() method
if isinstance(blob_data, bytes):
file_bytes = blob_data
else:
file_bytes = blob_data.content_as_bytes()
with open(local_file_path, "wb") as f:
f.write(file_bytes)
written = True
if written:
self.ledger_entry(action="read")
return written

def get_dataframe(
self, output="pandas", version="latest", pl_lazy: bool = False
) -> pd.DataFrame | pl.DataFrame:
Expand Down
157 changes: 157 additions & 0 deletions cfa/dataops/command.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,157 @@
"""
CLI tools for managing data operations in the CFA project.
"""

import os
from argparse import ArgumentParser

from rich.console import Console

from .catalog import datacat
from .utils import tree


def _get_dataset_namespaces() -> list[str]:
"""
Helper function to get a list of dataset namespaces.
"""
return datacat.__namespace_list__


def _get_stages_list(dataset_namespace: str) -> list[str]:
"""
Helper function to get a list of stages for a given dataset namespace.
"""
stages_start_with = ["load", "extract", "stage"]
if dataset_namespace not in _get_dataset_namespaces():
Console().print(
f"[bold red]Error:[/bold red] Dataset namespace '{dataset_namespace}' not found. Available datasets are:\n"
+ "\n".join(f"- {ds}" for ds in _get_dataset_namespaces())
)
return
dataset_dict = eval(f"datacat.{dataset_namespace}.__dict__")
stages = [
key
for key in dataset_dict.keys()
if any(key.startswith(prefix) for prefix in stages_start_with)
]
return sorted(stages)


def _get_versions_list(dataset_namespace: str, stage: str) -> list[str]:
"""
Helper function to get a list of versions for a given dataset namespace and stage.
"""
stages = _get_stages_list(dataset_namespace)
if stage not in stages:
Console().print(
f"[bold red]Error:[/bold red] Stage '{stage}' not found for dataset '{dataset_namespace}'. Available stages are:\n"
+ "\n".join(f"- {s}" for s in stages)
)
return
versions = eval(f"datacat.{dataset_namespace}.{stage}.get_versions()")
return versions


def get_available_data():
"""
Retrieve a list of available datasets for CFA.
"""
parser = ArgumentParser(description="Get list of available datasets")
parser.add_argument(
"-p", "--prefix", help="optional prefix filter", default=None
)
args = parser.parse_args()
datasets = datacat.__namespace_list__
if args.prefix:
datasets = [ds for ds in datasets if ds.startswith(args.prefix)]
formatted_list = "\n".join(f"- {dataset}" for dataset in sorted(datasets))
Console().print(f"[bold]Available Datasets:[/bold]\n{formatted_list}")


def get_dataset_stages():
"""
Retrieve stages of available datasets for CFA.
"""
parser = ArgumentParser(description="Get dataset stages")
parser.add_argument("dataset", help="full dataset namespace")
args = parser.parse_args()
dataset = args.dataset
stages = _get_stages_list(dataset)
formatted_stages = "\n".join(
f"- [red]{stage}[/red]" if stage == stages[-1] else f"- {stage}"
for stage in stages
)
disclaimer = "[italic yellow]Note: Stages in [red]red[/red] indicate the default stage for loading the dataset.[/italic yellow]"
Console().print(
f"[bold]Stages for {dataset}:[/bold]\n{formatted_stages}\n{disclaimer}"
)


def get_dataset_versions():
"""
Retrieve versions of available datasets for CFA.
"""
parser = ArgumentParser(description="Get dataset versions")
parser.add_argument("dataset", help="full dataset namespace")
parser.add_argument(
"--stage", "-s", help="specific stage to get version for", default=None
)
args = parser.parse_args()
dataset = args.dataset
available_stages = _get_stages_list(dataset)
if args.stage is None:
stage = available_stages[-1]
else:
stage = args.stage
versions = _get_versions_list(dataset, stage)
formatted_versions = "\n".join(
f"- [red]{version}[/red]" if version == versions[0] else f"- {version}"
for version in versions
)
Console().print(f"[bold]{dataset}[/bold]:\n{formatted_versions}")


def save_data_locally():
"""
Save a datasets to the local cache.
"""
parser = ArgumentParser(description="Download a dataset locally")
parser.add_argument("dataset", help="full dataset namespace")
parser.add_argument(
"location",
help="local location (directory) to save dataset to. If it does not exist, it will be created.",
)
parser.add_argument(
"--stage", "-s", help="specific stage to get version for", default=None
)
parser.add_argument(
"--version", "-v", help="specific version to get", default=None
)
parser.add_argument(
"--force", "-f", help="force re-download of data", action="store_true"
)
args = parser.parse_args()
dataset = args.dataset
stage = args.stage
version = args.version
if stage is None:
stages = _get_stages_list(dataset)
stage = stages[-1]
if version is None:
versions = _get_versions_list(dataset, stage)
version = versions[0]
local_path = os.path.abspath(args.location)
written = eval(
f"datacat.{dataset}.{stage}.download_version_to_local('{local_path}', version='{version}', force={args.force})"
)
if not written:
Console().print(
f"[bold yellow]Dataset '{dataset}' version '{version}' at stage '{stage}' is already present at location '{local_path}'. Use --force to re-download.[/bold yellow]"
)
return
else:
tree_output = tree(local_path, show_hidden=False)
Console().print(
f"[bold green]Dataset '{dataset}' version '{version}' at stage '{stage}' has been saved locally.[/bold green]\n\n{local_path}\n{tree_output}"
)
65 changes: 65 additions & 0 deletions cfa/dataops/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,8 @@
import glob
import os
from datetime import datetime
from itertools import islice
from pathlib import Path
from typing import Optional


Expand Down Expand Up @@ -124,3 +126,66 @@ def get_user() -> str:
return getpass.getuser()
except Exception:
return "unknown_user"


def tree(
dir_path: Path,
level: int = -1,
limit_to_directories: bool = False,
length_limit: int = 1000,
show_hidden: bool = False,
):
"""Given a directory Path object print a visual tree structure

Args:
dir_path (Path): the root directory path
level (int): how many levels deep to traverse, -1 for unlimited
limit_to_directories (bool): whether to only show directories
length_limit (int): maximum number of lines to print
show_hidden (bool): whether to show hidden files and directories
Returns:
str: visual tree structure
"""
space = " "
branch = "│ "
tee = "├── "
last = "└── "
dir_path = Path(dir_path) # accept string coercible to Path
files = 0
directories = 0

def inner(dir_path: Path, prefix: str = "", level=-1):
nonlocal files, directories
if not level:
return # 0, stop iterating
if limit_to_directories:
contents = [d for d in dir_path.iterdir() if d.is_dir()]
else:
contents = list(dir_path.iterdir())
pointers = [tee] * (len(contents) - 1) + [last]
for pointer, path in zip(pointers, contents):
if not show_hidden and path.name.startswith("."):
continue
if path.is_dir():
yield prefix + pointer + path.name
directories += 1
extension = branch if pointer == tee else space
yield from inner(
path, prefix=prefix + extension, level=level - 1
)
elif not limit_to_directories:
yield prefix + pointer + path.name
files += 1

lines = []
lines.append(dir_path.name)
iterator = inner(dir_path, level=level)
for line in islice(iterator, length_limit):
lines.append(line)
if next(iterator, None):
lines.append(f"... length_limit, {length_limit}, reached, counted:")
return (
"\n".join(lines)
+ f"\n{directories} directories"
+ (f", {files} files" if files else "")
)
18 changes: 18 additions & 0 deletions changelog.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,24 @@ The versioning pattern is `YYYY.MM.DD.micro(a/b/{none if release})

---

## [2025.10.31.0a]

### Added

- **New CLI commands** for dataset management:
- `dataops_datasets`: List available datasets with optional prefix filtering
- `dataops_stages`: View available stages for a specific dataset
- `dataops_versions`: List available versions for a dataset stage
- `dataops_save`: Download and save dataset versions locally
- **Tree utility function**: New `tree()` function in `utils.py` for displaying directory structures
- **Local data download**: `download_version_to_local()` method for BlobEndpoint to download datasets to local filesystem
- **Tests**: Comprehensive test coverage for new CLI commands and utilities
- New docs for CLI-tools

### Updated

- Documentation formatting: Changed "Description" to "Summary" with adjusted heading levels

## [2025.10.14.0a]

### Added
Expand Down
Loading