Data packages

Clustrix ships your function and its arguments. It does not ship your dataset. A function that opens "data/subjects.h5" finds that file on your laptop and does not find it on the worker, and no amount of decorating changes that. A data package is how a dataset travels: you name the files, clustrix moves them, and the function reads them back through the package on whichever machine it happens to be running on.

Nothing is inferred. A file moves because you named it, never because a string in your source code looked like a path. That restraint is deliberate: an upload triggered by the literal "s3://bucket/notes.log" is the worst failure mode available here, so declaration is the only route.

The shape of it

Three steps, and the middle one is the ordinary one:

# cluster-required: needs a configured cluster and a real data/ directory
import clustrix

# 1. Declare
subjects = clustrix.data_package("data/subjects.h5")

# 2. Pass it like any other argument
@clustrix.cluster(cores=8)
def fit(pkg):
    # 3. Dereference inside the function
    with open(pkg.path("subjects.h5"), "rb") as handle:
        return len(handle.read())

fit(subjects)

clustrix.data_package() accepts a path, a list of paths, a directory, or raw bytes. A directory expands to the files beneath it, and a list may mix files and directories freely. Whatever you name, the paths inside the package are relative to the common ancestor of everything named, so a function written against "data/subjects.h5" keeps working against "data/subjects.h5" on the worker. Pass base= when you want a different root.

Dereference with path(), which returns a filesystem path, or with read_bytes(), which returns the contents. Both check the file against the digest recorded when the package was built, so a file that arrived wrong is an error rather than a wrong answer. Both take no argument at all when the package holds exactly one file. A package dereferenced on the machine that built it reads your original files in place – no copy, no fetch – provided their digests still match what was packaged.

Here is the same round trip against the local backend, which needs no cluster and no network, and which this page executes for real every time it is checked:

import os

import clustrix
from clustrix import configure

configure(cluster_type="local")

os.makedirs("trials", exist_ok=True)
with open("trials/run1.csv", "w") as handle:
    handle.write("trial,rt\n1,0.42\n")

trials = clustrix.data_package("trials", force_local=True)
print(trials.filenames(), trials.total_bytes)   # ['run1.csv'] 16

@clustrix.cluster(cores=2)
def count_rows(pkg):
    with open(pkg.path("run1.csv")) as handle:
        return len(handle.readlines())

print(count_rows(trials))                       # 2

force_local=True says “carry the contents inside the object regardless of size”. Nothing is uploaded, and there is nothing to clean up afterwards.

Several packages travel as easily as one, and a list of them is walked the same way a single one is:

# cluster-required: needs a configured cluster to execute
import os

import clustrix

packages = [
    clustrix.data_package("data/subjects.h5"),
    clustrix.data_package("data/stimuli/"),
]

@clustrix.cluster(cores=4)
def summarize(pkgs):
    roots = clustrix.materialize_packages(pkgs)
    return [len(os.listdir(root)) for root in roots]

summarize(packages)

clustrix.materialize_packages() walks lists, tuples and dicts, writes every package it finds to local disk, and hands back the directory holding each. Anything that is not a package comes back unchanged.

Where the bytes actually live

Two places, and which one you get depends on size.

Small packages ride inside the object. Below stage_inline_max_bytes, the file contents are carried in the package itself, pickled alongside your function’s other arguments, and shipped over the transport that already moves the payload – SFTP for ssh and slurm, the payload channel for huggingface. No second transport, no remote store, nothing to delete.

Larger packages go to a private HuggingFace dataset repo. Above that threshold, the contents are uploaded and the package carries only the coordinates plus a digest per file. Three things follow, and you want to know all of them before it happens rather than after:

  1. It needs HuggingFace credentials. The token comes from hf_token in your config, then HF_TOKEN in the environment, then the cache that hf auth login writes. Without one, staging refuses and says so. The worker needs one too, from its own environment: clustrix deliberately does not pickle your token into the job payload, so a package staged remotely is unreadable on a worker with no HuggingFace credentials of its own.

  2. Clustrix creates a repo in your account. The first package that does not fit inline calls create_repo(private=True, exist_ok=True) for <namespace>/clustrix-data. The namespace is hf_namespace if you set one, otherwise hf_username, otherwise whatever the token’s whoami() reports. Set hf_data_repo to name a different repo outright.

  3. This applies on every backend. A slurm job with a package too big to inline still stages that package through HuggingFace, because that is the only remote store clustrix has.

Each package gets its own folder in that repo, keyed by a fresh identifier, so two packages never share a stored blob and deleting one cannot pull data out from under another. The cost of that is duplication: identical content packaged twice is stored twice.

Digests are computed locally, from your own files, and travel to the worker inside the function payload – which is an upload-only, local-origin artifact. Bytes fetched back out of the store are checked against those digests. In this way a tampered store is caught, because the expected digest never went through it.

The three size thresholds

Field

Default

Effect

stage_inline_max_bytes

1 MB

Under this, the package is inline. At or above it, the package is uploaded.

stage_warn_bytes

100 MB

At or above this, staging logs a warning before it starts. A transfer that takes minutes with no output is indistinguishable from a hang.

stage_max_bytes

5 GB

At or above this, staging raises clustrix.StagingError and names the largest file. Raise it if you genuinely mean to move that much over the network.

The inline threshold is measured on the serialized package, not on the raw data. Those two numbers are nowhere near each other once a package holds many small files, because every file also carries a relative path and a 64-character digest: ten thousand four-byte files are forty kilobytes of data and a 1.09 MB pickle. Measured on the data alone, that package would be a megabyte over the limit and still call itself inline, which is why the limit is applied to the pickle instead.

Deleting a package

Nothing is ever cleaned up for you. There is no TTL, no reaper, no eviction policy, and no deletion when a job finishes. cleanup_on_success governs the job directory on the cluster and does not touch staged data. Whether a dataset is still needed is your judgement rather than clustrix’s, so a package you staged stays staged, and stays billable, until you say otherwise.

The handle that created a package can end it:

# cluster-required: needs HuggingFace credentials
import clustrix

pkg = clustrix.data_package("data/subjects.h5")
pkg.delete()

delete() returns whether anything was actually removed. Calling it twice is not an error, and neither is calling it on a package somebody already cleaned up elsewhere. A remote copy that is present and refuses to delete does raise, because a warning there would leave you paying for storage you believe you released.

Four things it does not touch, each for its own reason:

  • The files you packaged. local_root points at your own data.

  • A directory you named yourself. If you called materialize(dest=...), that directory is yours; it may hold anything, and clustrix cannot tell what it put there from what was already there.

  • The cache directory above the package. Exactly one directory is removed, <local_cache_dir>/data-packages/<package id>, which clustrix created and which is keyed by an identifier nothing else uses.

  • The repo itself, only the package’s folder inside it. An account whose last package is deleted keeps an empty clustrix-data dataset, which is yours to remove by hand. A user who pointed hf_data_repo at a repo they own and care about would not thank clustrix for deleting it because the last package went away.

There is deliberately no context-manager form. A with block that quietly deleted the upload on the way out would be exactly the automatic cleanup this design rejects.

When the object is gone

Losing the handle does not mean losing the ability to clean up, because the object is not the only key to its own deletion:

# cluster-required: needs HuggingFace credentials
import clustrix

for record in clustrix.list_data_packages():
    print(record["package_id"], record["name"], record["total_bytes"])

clustrix.delete_data_package("0123456789abcdef0123456789abcdef")

clustrix.list_data_packages() returns one record per package in the store; an upload that was interrupted before its manifest went up appears with "complete": False. clustrix.delete_data_package() takes an id and validates it before anything reaches the Hub, since the id becomes a path in the store and an id that is not one is a deletion aimed somewhere else.

Keeping a package across sessions

The package object is the durable handle, and it is plain data: strings, ints, bytes. No client, no socket, no credential. So the way to keep a staged dataset across sessions is to pickle the object and load it later, and a copy loaded in a fresh interpreter still reaches the same remote data:

# cluster-required: needs HuggingFace credentials
import pickle

import clustrix

pkg = clustrix.data_package("data/subjects.h5")
with open("subjects.pkl", "wb") as handle:
    pickle.dump(pkg, handle)

# ... a week later, a different interpreter ...
with open("subjects.pkl", "rb") as handle:
    pkg = pickle.load(handle)

pkg.path("subjects.h5")   # still resolves
pkg.delete()              # still deletes

Credentials are deliberately not among the attributes that survive that round trip. A saved package must not be a token sitting on disk, so it re-authenticates from the ordinary config path every time it is used.

What it refuses

Every refusal below raises clustrix.StagingError before anything moves.

Paths that look like credentials. *.pem, *.key, .env, .netrc, id_rsa*, kubeconfig, anything under .ssh/, .aws/, .gnupg/, .kube/ or .config/gcloud/, a .git/config (which is one of the commonest places a personal access token ends up on disk), and so on. Both ends of a symlink are tested, since data/notes pointing at ~/.ssh/id_rsa is credential-shaped at the far end and innocuous at the near one. Pass allow_sensitive=True if you really do mean to move a keypair to the worker.

Anything that is not a regular file. Reading a fifo or /dev/zero does not fail, it blocks, so packaging one would hang with no output and no timeout. A refusal costs you one message.

A data repo that already exists and is public. create_repo(private=True, exist_ok=True) creates a private repo but does not make an existing public one private; it returns the repo as it is. Clustrix refuses rather than flipping the setting, because a repo can be public on purpose and silently changing someone’s visibility is its own incident. Make it private yourself, or point hf_data_repo somewhere else.

A package at or above stage_max_bytes – with the largest file named.

Two different files that would land on the same name. One would silently overwrite the other on the worker and the run would produce a wrong answer rather than an error. Naming the same file twice is harmless and collapses.

Files that change while they are being staged. Files go up by path, so the bytes on the wire are whatever the file held at upload time rather than what was hashed a moment earlier. Clustrix re-hashes after the commit and, if anything moved, removes the folder it just created and tells you which files changed.

When not to use this

If the data is already reachable from the worker, staging it is pure waste. On an HPC cluster with a shared filesystem your /scratch directory is visible from every compute node, so the file your function opens is already there and moving a copy of it through HuggingFace buys nothing but a transfer and a storage bill. The same goes for data in object storage the worker can authenticate to.

Use the read-only filesystem utilities to check. cluster_exists answers the question directly:

from clustrix import cluster_exists
from clustrix.config import ClusterConfig

config = ClusterConfig(cluster_type="local", local_work_dir=".")

if cluster_exists("setup.py", config):
    print("already there; nothing to stage")

Data packages are for the case where that answer is no.

API

clustrix.data_package(source, *, name=None, base=None, config=None, force_local=False, allow_sensitive=False, filename='data.bin')[source]

Package data or files into an object a @cluster function can read.

Parameters:
  • source (Union[str, Path, bytes, Sequence[Union[str, Path]]]) – A path, a list of paths, or raw bytes. Directories expand to the files beneath them.

  • name (Optional[str]) – Label for messages; defaults to the source’s basename.

  • base (Union[str, Path, None]) – What the package’s relative paths are relative to. Defaults to the common ancestor of everything named, which is what lets a function keep using the paths it already uses.

  • config – A ClusterConfig; the global one by default.

  • force_local (bool) – Carry the contents inside the object regardless of size. Nothing is uploaded and there is nothing to clean up.

  • allow_sensitive (bool) – Permit paths that look like credentials.

  • filename (str) – The name raw bytes get inside the package.

Return type:

DataPackage

Returns:

A DataPackage. Pass it – or a list of them – to a @cluster-decorated function as an ordinary argument.

Raises:

StagingError – on a missing file, a credential-shaped path, a path that escapes the package root, or a package at or above stage_max_bytes.

class clustrix.DataPackage(name, package_id, files, local_root=None, repo_id=None, path_in_repo=None, inline=None, _materialised=None)[source]

A named bundle of files that can travel to a worker and be read there.

Build one with data_package(); the constructor here is the low-level form and does no staging of its own.

name

A label, used in messages and in the materialisation directory.

package_id

Unique per package. Also the folder name in the remote store, which is what keeps two packages from sharing a blob.

files

One PackagedFile per file, in declaration order.

local_root

Where the files are on the machine that built the package. None for a package built from in-memory data. Used as a zero-cost source when the same machine dereferences the package.

repo_id / path_in_repo

The private remote location, or None for an inline package.

inline

relpath -> bytes when the payload rides inside the object.

Every attribute is plain data – strings, ints, bytes. No HF client, no socket, no file handle, so the object pickles cleanly and a copy loaded in a fresh interpreter still reaches the same remote data. Credentials are deliberately not among those attributes: a saved package must not be a token sitting on disk, so it re-authenticates from the ordinary config path whenever it is used.

property is_inline: bool

Whether the bytes ride inside this object rather than in the store.

filenames()[source]

The package-relative paths, in declaration order.

Return type:

List[str]

path(relpath=None, config=None)[source]

Local filesystem path to one file, fetching it if it is not here yet.

Called with no argument on a single-file package. This is the “on demand” half: nothing is fetched until something asks for it, and a package dereferenced on the machine that built it reads the original files without copying them.

Return type:

str

read_bytes(relpath=None, config=None)[source]

The contents of one file, verified against its recorded digest.

Return type:

bytes

materialize(dest=None, config=None)[source]

Put every file on local disk and return the directory holding them.

Files keep the relative paths they had when the package was built, so a function that opened "data/x.csv" opens "data/x.csv" under the returned root on either machine.

Return type:

str

exists(config=None)[source]

Whether the remote copy is still there. Inline packages: always.

Return type:

bool

__getstate__()[source]

Pickle without the materialisation path.

Where this package was last unpacked is true of one machine at one moment. Carrying it into a pickle means a worker – or this machine a week later – can find a stale directory at that path and reuse it. Everything else here is durable; this one field is not.

Return type:

Dict[str, Any]

delete(config=None)[source]

Remove the remote copy and any clustrix-owned local cache of it.

Returns whether anything was actually removed remotely. Calling this twice, or on a package that was already cleaned up elsewhere, is not an error – but a remote copy that is present and cannot be deleted raises, because a warning here would leave the caller paying for storage they believe they released.

Nothing calls this for you. There is no TTL and no reaper: whether a dataset is still needed is the user’s judgement, not clustrix’s.

The files you packaged are never touched. local_root points at the user’s own data; deleting that because a transfer was cleaned up would be indefensible. Only the remote folder, and the copy clustrix wrote into its own cache, are removed.

A directory you named yourself is never touched either. If you called materialize(dest=...), that directory is yours – it may hold anything, and clustrix has no way to know what it created there versus what was already in it. It used to be removed recursively, which deleted whatever else the caller kept alongside the data. Clear it yourself if you want it gone.

The repo itself is never touched either, only the package’s folder inside it. hf_data_repo may well point at a repo the user owns and cares about. An account whose last package is deleted keeps an empty clustrix-data dataset, which is theirs to remove.

Return type:

bool

__init__(name, package_id, files, local_root=None, repo_id=None, path_in_repo=None, inline=None, _materialised=None)
clustrix.materialize_packages(value, config=None)[source]

Walk a structure and materialise every DataPackage in it.

Convenience for a worker that would rather have paths than handles:

roots = clustrix.materialize_packages(packages)

Lists, tuples, and dicts are walked; anything else is returned unchanged.

Return type:

Any

clustrix.list_data_packages(config=None)[source]

Every package clustrix has staged in the remote store.

A user who lost the handle still needs a way to see what is costing them storage, so the object is not the only key to its own deletion. Returns a list of manifests; incomplete uploads appear with "complete": False.

Return type:

List[Dict[str, Any]]

clustrix.delete_data_package(package_id, config=None)[source]

Delete a staged package by id, without needing the object.

The counterpart to list_data_packages(). Returns whether anything was removed; a package that is already gone is not an error, a package that is there and will not delete raises.

The id is validated first, before anything reaches the Hub. It becomes a path in the store, so an id that is not one is a deletion aimed somewhere else: "" addressed the whole packages/ prefix and removed every package in the account, and "../README.md" climbed out of the prefix and removed a file that was never a package.

Return type:

bool

exception clustrix.StagingError[source]

Anything that stops a package being built, moved, or removed.