Skip to content

Latest commit

 

History

History
307 lines (224 loc) · 21.6 KB

File metadata and controls

307 lines (224 loc) · 21.6 KB

CLAUDE.md

This file provides guidance to Claude Code (claude.ai/code) when working with code in this repository.

Overview

Clustrix is a Python distributed computing framework that runs Python functions on remote compute using a @cluster decorator.

Backends, and how far each is actually proven — keep this honest, it is the first thing anyone reads:

cluster_type Status
slurm Verified end to end against a real scheduler
ssh Verified end to end against a real GPU host
huggingface Verified end to end against real HF Jobs containers
local Runs in-process via local_executor.py

That is the whole list — clustrix.config.SUPPORTED_CLUSTER_TYPES. Evidence for the verified rows is regenerated by scripts/verify_cluster_usecases.py and committed under docs/evidence/.

Backends that are NOT supported. pbs, sge, kubernetes and the AWS / GCP / Azure / Lambda Cloud VM providers are absent, and naming one in cluster_type raises a ValueError. So are the cost monitoring and cloud pricing API (cost_tracking_decorator, get_cost_monitor, start_cost_monitoring, generate_cost_report, get_pricing_info) and the HuggingFace Spaces provider — which is a different thing from cluster_type="huggingface" (HuggingFace Jobs), and that one is supported. Each absent backend is planned for a future update and has a tracking issue; do not document any of them as working:

Not supported Issue
PBS #140
SGE #141
Kubernetes #142
AWS #143
GCP #144
Azure #145
Lambda Cloud #146

scripts/aws/ stays: it is operator cleanup tooling, not an execution backend.

Development Commands

Setup

# Install package with development dependencies
pip install -e ".[dev]"

# Install with the Jupyter widget
pip install -e ".[widget,dev]"

Code Quality

CRITICAL: Always run these checks before committing or pushing code:

# Use the comprehensive quality check script (recommended)
python scripts/check_quality.py

# Or run individual checks:
black clustrix/          # Format code
flake8 clustrix/         # Run linter  
mypy clustrix/           # Type checking
pytest --cov=clustrix    # Run tests with coverage

Pre-commit hooks are installed - they will automatically run black, flake8, and mypy before each commit. If any check fails, the commit will be blocked until fixed.

Never push without running quality checks - GitHub Actions will fail if code doesn't pass black, flake8, and mypy.

Architecture

Core Components

  1. @cluster Decorator (clustrix/decorator.py): Main user interface for marking functions for remote execution. Supports resource specification and automatic loop parallelization.

    Source availability: serialization itself does not need the function's source — serialize_function/deserialize_function round-trip a function created by exec() and return the correct answer, because dill and cloudpickle work from the code object. The one thing that needs source is AST loop parallelization, which calls inspect.getsource(); when the source is unavailable that step is skipped and the function ships as-is. There is no complexity analysis and no function flattening — those were deleted (#89/#90), so do not document them. Do not describe any of this as "functions cannot be serialized in the REPL" either; that claim is false and led to a fabricated-result bug.

    cores does nothing on the local path. @cluster(cores=8) with no cluster_host runs the function in the calling process, sequentially. decorator.py's local branch is a bare return func(*args, **func_kwargs) whenever no parallelizable loop is found, which is nearly always. This is issue #152, it is a defect rather than a design position, and docs/source/notebooks/local_parallel_comparison.ipynb measures it. Real local parallelism lives in LocalExecutor (local_executor.py), which picks ProcessPoolExecutor or ThreadPoolExecutor via choose_executor_type (local_executor.py:339): pickle-test the function and every argument (failure → threads), then substring-scan the source for I/O markers (open(, requests., time.sleep, …) (hit → threads), else processes. use_threads overrides.

    cluster_type is not a @cluster keyword. @cluster(cores=8, cluster_type="local") logs @cluster received unrecognised option(s) on every call and is ignored. The backend is set with configure(cluster_type=...). The only extra keywords the decorator accepts are hf_token, hf_username, hf_flavor, hf_timeout, hf_namespace and key_file; check every example you write for this mistake.

    auto_gpu_parallel does nothing. It is accepted, it warns, and there is no automatic cross-GPU parallelization to select. Parallelize across GPUs inside the function.

  2. ClusterExecutor (clustrix/executor_core.py): Central execution engine. Note that clustrix/executor.py is a 39-line backward-compatibility shim that re-exports it; the implementation is split across:

    • executor_core.py — the ClusterExecutor class, dispatch, result retrieval and verification
    • executor_connections.py — SSH/SFTP connection management via Paramiko
    • executor_schedulers.py — SLURM and SSH submission
    • executor_scheduler_status.py — scheduler status polling and error extraction
    • hf_jobs.py — the HuggingFace Jobs backend (HFJobsManager)
  3. Configuration System (clustrix/config.py): Singleton configuration management supporting:

    • YAML/JSON file loading
    • Hierarchical configuration (defaults → file → runtime)
    • Standard location discovery (~/.clustrix/, /etc/clustrix/)

Data staging. Clustrix does not move your data. The pickled function and its pickled arguments travel; nothing else does. clustrix/file_packaging.py and clustrix/dependency_analysis.py can build a ZIP of a function's local dependencies, and nothing on the execution path calls either of them — confirm with grep -rn "FilePackager\|package_function" clustrix/decorator.py clustrix/executor_core.py clustrix/executor_connections.py clustrix/utils.py, which returns empty. The cluster_* filesystem helpers are read-only: no cluster_put, no cluster_get. Declared data goes through clustrix/staging.py (data_package, DataPackage, materialize_packages, list_data_packages, delete_data_package, StagingError) — declaration only, never inference from source. Do not document any inferred staging as a feature; there is none, deliberately. Two facts about it that are easy to get wrong: a package too big to inline is uploaded to a private HuggingFace dataset repo clustrix creates in the user's account, on every backend, and nothing is ever deleted automatically — no TTL, no reaper, and cleanup_on_success does not touch staged data. The user-facing guide is docs/source/data_packages.rst.

  1. Utilities (clustrix/utils.py):

    • Function serialization using cloudpickle/dill
    • Environment capture and replication
    • AST-based loop detection for parallelization
    • Cluster-specific job script generation
  2. Filesystem Utilities (clustrix/filesystem.py): Unified filesystem operations for local and remote clusters:

    • ClusterFilesystem class: Core implementation handling local and SSH-based remote operations
    • Convenience functions: cluster_ls(), cluster_find(), cluster_stat(), cluster_exists(), cluster_isdir(), cluster_isfile(), cluster_glob(), cluster_du(), cluster_count_files()
    • Data classes: FileInfo for file metadata, DiskUsage for directory usage statistics
    • Automatic connection management: SSH connections handled transparently
    • Config-driven: Uses ClusterConfig to determine local vs remote operations

Key Design Patterns

  • Strategy Pattern: Different cluster types handled via specific submission methods
  • Decorator Pattern: Clean API through @cluster decorator
  • Singleton Pattern: Global configuration instance
  • Template Pattern: Job script generation with cluster-specific variations

Workflow

  1. User decorates function with @cluster
  2. Function and dependencies serialized with cloudpickle
  3. Job submitted to cluster:
    • Remote directory created
    • Serialized data uploaded
    • Python environment setup (venv/conda)
    • Job script generated and submitted
  4. Remote execution with unpickling
  5. Results polled and downloaded

Two-venv remote execution

Scheduler and SSH backends do not run the function in one interpreter. generate_two_venv_execution_commands (utils.py) emits a bash script containing three python -c "..." programs:

  1. VENV1 deserialize — read function_data.pkl, write function_deserialized.pkl
  2. VENV2 execute — run the function, write result_raw.pkl
  3. VENV1 serialize — write result.pkl and its HMAC tag

VENV1 holds clustrix's own serialization dependencies; VENV2 holds the user's replicated environment. Every handoff must use dill/cloudpickle, never stdlib pickle — pickle serializes a function by qualified name, which cannot be resolved in a fresh interpreter, and that asymmetry is what made remote execution fail for every __main__ function. Any change here must keep each serialize/deserialize pair symmetric; tests/unit/test_two_venv_execution.py enforces this.

Important Considerations

  • The project is beta. Version strings live in pyproject.toml, setup.py, clustrix/__init__.py and docs/source/conf.py and must be kept identical.
  • SSH-based clusters require proper key setup or password authentication
  • Host keys are verified by default. Every paramiko connection goes through clustrix/ssh_security.py::configure_host_key_policy. Never call set_missing_host_key_policy(paramiko.AutoAddPolicy()) directly — the opt-out is ClusterConfig.ssh_host_key_policy="auto_add", and it is honoured only from a trusted configuration source: host_key_policy_name downgrades auto_add to reject when config_source_is_trusted says no, because a weakening of verification is a security decision and it persists in known_hosts. hf_image follows the same rule, for the same reason — a staged HF job hands CLUSTRIX_HF_TOKEN to whatever image it names. On that opt-out clustrix creates ~/.ssh/known_hosts when it is absent (0700 directory, 0600 file) and loads it with load_host_keys, because paramiko's AutoAddPolicy persists a key only when _host_keys_filename is set — without both steps auto_add accepts every host and saves none. The path has exactly one definition, ssh_security.user_known_hosts_path; ssh_utils imports it rather than recomputing it.
  • Remote environments are recreated from the local environment's freeze output
  • Job scripts are bash-based with scheduler-specific directives
  • Every stored credential is released through one function. clustrix/credential_release.py::release_credential(target, ...) is the only place a stored SSH secret is handed out, and its first positional parameter is the recipient — a frozen CredentialTarget that names the hostname, the username and who chose the hostname. FlexibleCredentialManager._ensure_credential_unchecked raises for any caller that is not that module, always, in production. Never obtain a secret any other way: issue #167 was reported as one leak and closed as seven, and seven call sites for one decision is a decision with no home. Provenance is an argument (ClusterConfig.from_file_content(mapping, source)), not ambient context. tests/unit/test_every_credential_goes_through_one_gate.py fails when a new secret-bearing structure appears anywhere in clustrix/.
  • Results are HMAC-verified before they are deserialized. Loading a pickle executes code, so a file fetched from a remote host is a remote-to-local code-execution path. result.pkl is signed with a per-job key and checked by executor_core.py before dill.loads. Any new remote-origin byte stream the caller parses must be authenticated the same way.
  • Automatic cleanup of remote files configurable via cleanup_on_success

Common Tasks

Adding New Cluster Type Support

There is no ClusterType enum — ClusterConfig.cluster_type is a plain str. The supported values are local, ssh, slurm, huggingface, declared once in clustrix.config.SUPPORTED_CLUSTER_TYPES.

  1. Add the value to SUPPORTED_CLUSTER_TYPES; the CLI's click.Choice and both notebook widgets' dropdowns read that tuple, so they cannot drift apart. Verify rather than trust this sentence — grep -rn 'list(SUPPORTED_CLUSTER_TYPES)' clustrix/ must show three call sites (cli.py, modern_notebook_widget.py, notebook_magic_widget.py). Until #165 it showed two: notebook_magic_widget.py spelled the four values out, so that menu could drift, and this line claimed otherwise.
  2. Implement submission in the appropriate executor_*.py module and dispatch from ClusterExecutor in executor_core.py
  3. Add status checking to get_job_status / executor_scheduler_status.py
  4. Update job script generation in utils.py if needed — reuse job_execution_lines() rather than writing another variant
  5. Do not mark it supported until a real job has run on real hardware and the evidence is committed. That gate is why pbs, sge, kubernetes and the cloud VM providers are absent (#140–#146).

Using Filesystem Utilities

from clustrix import cluster_ls, cluster_find, cluster_stat, cluster_exists, cluster_glob
from clustrix.config import ClusterConfig

# Configure for local or remote operations
config = ClusterConfig(
    cluster_type="slurm",  # or "local" for local operations
    cluster_host="cluster.example.edu",
    username="researcher",
    remote_work_dir="/scratch/project"
)

# List directory contents (works locally and remotely)
files = cluster_ls("data/", config)

# Find files by pattern
csv_files = cluster_find("*.csv", "datasets/", config)

# Check file existence
if cluster_exists("results/output.json", config):
    print("Results already computed!")

# Get file information
file_info = cluster_stat("large_dataset.h5", config)
print(f"Dataset size: {file_info.size / 1e9:.1f} GB")

# Use with @cluster decorator for data-driven workflows
@cluster(cores=8)
def process_datasets(config):
    # Find all data files on the cluster
    data_files = cluster_glob("*.csv", "input/", config)
    
    results = []
    # Sequential: auto-parallelization needs a literal range() and a callee
    # that accepts the chunk keywords.
    for filename in data_files:
        # Check file size before processing
        file_info = cluster_stat(filename, config)
        if file_info.size > 100_000_000:  # Large files
            result = process_large_file(filename, config)
        else:
            result = process_small_file(filename, config)
        results.append(result)
    
    return results

Debugging Remote Execution

  • Check remote logs in ~/.clustrix/jobs/{job_id}/
  • Enable SSH debug mode in Paramiko connection
  • Verify environment setup with pip freeze comparison
  • Check scheduler-specific logs (SLURM: slurm-*.out)

Configuration Priority

  1. Runtime parameters — configure(...) and @cluster(...) keywords (highest priority)
  2. Configuration file (clustrix.yml, discovered in ~/.clustrix/ and /etc/clustrix/)
  3. Default values (lowest priority)

There is no environment-variable level. Nothing reads a CLUSTRIX_<FIELD> variable; grep for it before believing otherwise. Only two environment variables are consulted at all, and neither sets a config field:

  • CLUSTRIX_CONFIG_DIR — where configuration files are looked for and saved
  • whatever ClusterConfig.password_env_var names — read by the auth fallback to supply a password, and only a password

save_to_file omits secret-bearing fields by default, so password_env_var is the only supported channel for getting a credential in without writing it to disk. A general environment-variable overlay would be a reasonable feature; it has not been built, and the documentation must not imply it has.

⚠️ MANDATORY PRE-COMMIT WORKFLOW ⚠️

EVERY SINGLE TIME before committing/pushing:

  1. Run quality checks repeatedly: python scripts/pre_push_check.py
    • This runs black (formatting), flake8 (linting), mypy (type checking), and pytest (tests)
    • It automatically retries up to 5 times until ALL checks pass
    • Black may auto-fix formatting issues, so subsequent runs check if everything is clean
  2. Only commit and push when ALL checks pass

CRITICAL: You must run checks repeatedly until they ALL pass. Don't just run once!

Pre-commit hooks will block commits that fail quality checks.

The GitHub Actions CI will fail if code doesn't pass black, flake8, mypy, and pytest. Always run the full cycle locally first.

Testing Guidelines

The mocking policy, stated once

There is exactly one mocking policy, and it is this one. Do not cite .claude/CLAUDE.md's "do not use mock services for anything ever" against it, or any wording about mocking external dependencies; this list governs.

  1. Real first, always. A capability may not be marked working until it has been exercised against the real thing — a real cluster, a real API, a real file on disk, a real socket. A test that has only ever passed against a mock is evidence of nothing.
  2. Mocks are a cost-control measure, never a correctness argument. Once a real call has verified the contract, a mocked test using the same call syntax may stand in for it in CI to avoid per-run API fees and credential requirements. Re-verify against the real service when the contract could have changed.
  3. A mock may never be a fallback. If real functionality is unavailable, the test must fail or raise. Silently substituting a mock turns a broken feature into a green test.
  4. Production code must never know it is being tested. No isinstance(x, Mock), no test-only branches, no importable module of fake widgets. This is issue #116; grep -rn "unittest.mock\|MagicMock\|isinstance(.*Mock" clustrix/ must stay empty.
  5. Never weaken a test to make it pass. If a test fails, fix the code. If the test itself asserts wrong behaviour, say so explicitly and rewrite the assertion — do not quietly relax it.

20 of the 152 test modules use unittest.mock in ways that violate (1) and (2); replacing them is issue #117. New tests must not add to that number. Recount with grep -lE "unittest\.mock|Mock\(|MagicMock\(|@patch" $(find tests -name "test_*.py") | wc -l before quoting a figure.

Test Organization

Unit Tests (tests/ excluding real_world/ and integration/):

  • Run in GitHub Actions CI
  • No external resources required; fast
  • Included in coverage reports

Real-World Tests (tests/real_world/):

  • Test actual cluster functionality, API calls, SSH connections
  • tests/real_world/conftest.py applies @pytest.mark.real_world to every item in that directory automatically. Do not rely on a per-file decorator; a file that lacks one is still marked, and a file that has one adds nothing.
  • Run manually via the real-world-tests workflow (workflow_dispatch), which is gated on the required secrets being present

Integration Tests (tests/integration/):

  • These provision real, billable AWS resources. They refuse to run unless CLUSTRIX_ALLOW_BILLABLE=1 is set.
  • The guard reads config.args, not config.invocation_params.args. This is deliberate: invocation_params.args misses cases that red-teaming the guard turned up. Do not "simplify" it.

Pre-Push Hook Workflow

A custom pre-push hook (.git/hooks/pre-push) automatically runs real-world tests when:

  1. Cluster credentials are available (environment variables, etc.)
  2. Real-world test runner script exists

Hook behavior:

  • ✅ With credentials: Runs all real-world tests and blocks push on failure
  • ⚠️ Without credentials: Skips with warning, allows push
  • 🔧 Script missing: Falls back to basic pytest real-world test run

Running Tests Manually

# Everything that is safe to run without credentials or money.
# The --ignore flags are belt-and-braces alongside the marker: this is
# exactly what CI runs, so a green run here means a green run there.
pytest tests/ -m "not real_world" --ignore=tests/real_world --ignore=tests/integration

# Real-world tests only -- makes real SSH and cloud API calls
pytest tests/real_world/ -m real_world

# Using real-world test runner (when available)
python scripts/run_real_world_tests.py --filesystem
python scripts/run_real_world_tests.py --ssh  
python scripts/run_real_world_tests.py --api
python scripts/run_real_world_tests.py --visual

External Function Validation

  • External Function Validation: When testing external functions and features (i.e., anything that we cannot directly test locally), we MUST validate that those functions work (without resorting to local fallbacks) at least once before we can check the issue off as completed. Use GitHub issue spec criteria and comments to track what has been validated. We can store API keys, usernames, and/or passwords locally (if we can do so securely) in order to enable this. Once we have verified functionality (again, WITHOUT resorting to fallbacks), we can check off that functionality as "tested" and then use mocked functions or objects in pytests to avoid incurring excessive API fees.