Metadata-Version: 2.5
Name: lgopy
Version: 2.0.0
Summary: LgoPy is an open-source library for building multimodal data science pipelines.
Project-URL: Homepage, https://github.com/OpenSciML/lgopy
Project-URL: Repository, https://github.com/OpenSciML/lgopy
Project-URL: Documentation, https://opensciml.github.io/lgopy/
Project-URL: Issues, https://github.com/OpenSciML/lgopy/issues
License-Expression: Apache-2.0
License-File: LICENSE
Keywords: blocks,data-science,phenomics,pipelines,workflow
Classifier: Development Status :: 3 - Alpha
Classifier: Intended Audience :: Developers
Classifier: Intended Audience :: Science/Research
Classifier: Programming Language :: Python :: 3
Classifier: Topic :: Scientific/Engineering
Classifier: Topic :: Software Development :: Libraries :: Python Modules
Requires-Python: >=3.10
Requires-Dist: fsspec>=2024.10.0
Requires-Dist: lgopy-catalog[rag]==2.0.0
Requires-Dist: multimethod>=1.12
Requires-Dist: pydantic>=2.8.2
Requires-Dist: python-dotenv>=1.1.1
Requires-Dist: rich>=13.8.0
Requires-Dist: scikit-learn>=1.5.8
Provides-Extra: catalog
Requires-Dist: lgopy-catalog==2.0.0; extra == 'catalog'
Provides-Extra: catalog-rag
Requires-Dist: lgopy-catalog[rag]==2.0.0; extra == 'catalog-rag'
Provides-Extra: dask
Requires-Dist: dask[complete]>=2026.3.0; extra == 'dask'
Provides-Extra: dev
Requires-Dist: autoflake>=2.3.1; extra == 'dev'
Requires-Dist: bandit>=1.7.9; extra == 'dev'
Requires-Dist: black>=24.8.0; extra == 'dev'
Requires-Dist: isort>=6.0.1; extra == 'dev'
Requires-Dist: keyrings-google-artifactregistry-auth>=1.1.2; extra == 'dev'
Requires-Dist: packaging>=24.0; extra == 'dev'
Requires-Dist: pylint>=3.2.7; extra == 'dev'
Requires-Dist: pytest>=8.0; extra == 'dev'
Requires-Dist: twine>=5.1.1; extra == 'dev'
Provides-Extra: docs
Requires-Dist: mkdocs-material>=9.5.40; extra == 'docs'
Requires-Dist: mkdocs>=1.6.1; extra == 'docs'
Provides-Extra: examples
Requires-Dist: matplotlib>=3.9.2; extra == 'examples'
Requires-Dist: pandas>=2.2.3; extra == 'examples'
Requires-Dist: requests>=2.32.3; extra == 'examples'
Requires-Dist: xarray>=2024.10.0; extra == 'examples'
Provides-Extra: vector-hub
Requires-Dist: faiss-cpu>=1.13.2; extra == 'vector-hub'
Requires-Dist: langchain-community>=0.4.1; extra == 'vector-hub'
Requires-Dist: langchain-huggingface>=1.2.1; extra == 'vector-hub'
Requires-Dist: langchain>=1.2.13; extra == 'vector-hub'
Requires-Dist: sentence-transformers>=5.3.0; extra == 'vector-hub'
Description-Content-Type: text/markdown

# LgoPy

**Build data-processing pipelines from reusable Python blocks.**

LgoPy turns data transformations, model inference, and analysis operations into
configurable blocks that can be tested independently and composed into pipelines.
Blocks follow scikit-learn's fit/transform conventions, while data adapters
control how they process individual records, arrays, or complete datasets.

Shared stores capture metadata and artifacts throughout a run. With
`lgopy-catalog`, you can package and publish versioned blocks, discover them
through metadata filters or semantic search, and load them into new pipelines.
Develop blocks in standalone Python projects, save pipeline definitions as JSON,
and integrate the same processing logic into your applications.

[Install](#installation) · [Quickstart](#quickstart) ·
[Features](#feature-overview) · [Documentation](#documentation) ·
[Contributing](#development-and-contributing)

## Feature overview

| Capability | What you can do |
| --- | --- |
| **Composable blocks** | Define typed `call()` operations, configure constructor parameters, and initialize resources through `setup()`. |
| **Pipeline validation** | Check adjacent block signatures and validate serialized step definitions before execution. |
| **Data-processing adapters** | Register dispatch rules for custom containers and select item-level or dataset-level processing from block metadata. |
| **Metadata and artifacts** | Share context between steps, attach custom attributes, and save outputs in memory or through fsspec. |
| **Schemas and serialization** | Generate constructor/input/output schemas, serialize block configuration, and save pipelines as JSON. |
| **Versioned registries** | Register blocks using decorators or class attributes and instantiate a specific version. |
| **Portable packages** | Build folders or ZIP archives containing source, dependency requirements, schemas, and manifests. |
| **Block catalogs** | Publish, list, filter, materialize, load, instantiate, and remove stored block packages. |
| **Semantic discovery** | Search indexed packages by meaning using Gemini, Ollama, or custom embedding and vector-index adapters. |
| **Execution hooks** | Observe pipeline start, step start/completion, success, and failure through callbacks. |
| **Extensible file IO** | Register readers and exporters by file extension for application-specific data formats. |
| **Prompt assembly** | Experiment with pipeline selection from registered blocks using keyword matching or a supplied language model. |

## Installation

Requires **Python 3.10 or newer**:

```bash
pip install lgopy
```

The current core dependency list includes `lgopy-catalog[rag]`. Installing the
package does not start PostgreSQL or an embedding service; configure those only
when using semantic search. Additional extras support examples and development:

```bash
pip install "lgopy[examples]"  # pandas, xarray, plotting, and example clients
pip install "lgopy[docs]"      # MkDocs and Material theme
pip install "lgopy[dev]"       # pytest, package checks, and source-processing tools
```

Install the dependencies required by your own blocks as well. For remote fsspec
storage, install the corresponding backend, such as `gcsfs` for `gs://` URLs, and
configure credentials. See the [installation guide](mkdocs/installation.md) for
all extras and their current limitations.

## Quickstart

Save this as `normalize_block.py`. Later examples reuse this block:

```python
from lgopy.core import Block, BlockHub, LgoPipeline

@BlockHub.register(
    name="normalize",
    version="1.0.0",
    display_name="Normalize measurements",
    category="numeric",
    description="Divide each measurement by a configurable scale.",
    tags=["numeric", "normalization"],
    extras={"transform_scope": "dataset_item", "batch_independent": True},
)
class Normalize(Block):
    """Normalize individual measurements.

    Args:
        scale: Nonzero divisor applied to each input value.
    """

    def __init__(self, scale: float = 100.0) -> None:
        """Configure normalization.

        Args:
            scale: Nonzero divisor applied to each input value.
        """
        super().__init__()
        self.scale: float = scale

    def call(self, value: float) -> float:
        """Normalize one measurement.

        Args:
            value: Measurement to normalize.

        Returns:
            Input divided by the configured scale.

        Raises:
            ZeroDivisionError: If scale is zero.
        """
        return value / self.scale

if __name__ == "__main__":
    pipeline = LgoPipeline.from_steps(Normalize(scale=10.0))
    assert pipeline([10.0, 20.0, 30.0]) == [1.0, 2.0, 3.0]
```

Run `python normalize_block.py`. The built-in list adapter calls the block once
per item. NumPy arrays are passed to `call()` whole, allowing array operations
when the block supports them.

Use `setup()` to initialize models or other resources during fitting.
`pipeline(data)` performs fitting and transformation; repeated calls can run
setup again. Calling `block(data)` directly transforms without fitting.

## Configure, validate, and restore pipelines

Importing the quickstart module registers `normalize` in the current process:

```python
from normalize_block import Normalize
from lgopy.core import BlockHub, LgoPipeline

block = BlockHub.create("normalize", version="1.0.0", scale=10.0)
assert block.to_dict() == {"scale": 10.0}
print(Normalize.schema())

steps = [
    {"block": "normalize", "version": "1.0.0", "args": {"scale": 10.0}},
]
report = LgoPipeline.validate(steps)
if not report["valid"]:
    raise ValueError(report["issues"])

pipeline = LgoPipeline.from_list(steps)
pipeline.save("pipeline.json")
restored = LgoPipeline.from_file("pipeline.json")
assert restored([10.0, 20.0]) == [1.0, 2.0]
```

Use `to_json()` and `from_json()` for JSON strings. Saved definitions record block
names, versions, and constructor arguments; they do not bundle input data,
dependencies, artifacts, or fitted model state. Validation inspects block
signatures and configuration, so also test representative inputs.

Constructor annotations generate configuration schemas; `typing.Annotated` can
provide field descriptions. Keep input and output annotations accurate so schema
discovery and compatibility checks reflect the block's actual behavior.

The decorator is optional: declare metadata such as `name`, `version`, `tags`,
and `extras` as class attributes, then call `BlockHub.register(MyBlock)` when
registry lookup is needed. Avoid conflicting declarations between the two styles.
See [registries and catalogs](mkdocs/guide/registries-catalogs.md).

## Control data processing with adapters

Register `apply_transform` for a custom container to decide how each block
receives its data. For example, this adapter reads the scope declared in the
block's `extras` metadata:

```python
from dataclasses import dataclass
from typing import Any

from lgopy.core import Block, LgoPipeline, apply_transform
from normalize_block import Normalize

@dataclass
class Measurements:
    """Container of numeric values for adapter-controlled processing."""

    values: list[float]

@apply_transform.register
def apply_measurements(data: Measurements, block: Block) -> Any:
    """Dispatch using the block's declared scope.

    Args:
        data: Complete container supplied to this pipeline step.
        block: Block declaring dataset_item or dataset scope in extras.

    Returns:
        A container of item outputs, or the dataset-level block's result.

    Raises:
        ValueError: If the block declares an unsupported scope.
    """
    scope = getattr(block, "extras", {}).get("transform_scope", "dataset_item")
    if scope == "dataset":
        return block.call(data)
    if scope == "dataset_item":
        return Measurements([block.call(value) for value in data.values])
    raise ValueError(f"Unsupported transform_scope: {scope!r}")

pipeline = LgoPipeline.from_steps(Normalize(scale=10.0))
assert pipeline(Measurements([10.0, 20.0])).values == [1.0, 2.0]
```

The adapter implements the scope convention; flags alone do not change LgoPy's
built-in dispatch. Keep adapters in an importable module and load it in every
execution process.

For large inputs, an application can run the pipeline on bounded batches of
independent items. Check scope and batch independence before splitting the data:
dataset-wide statistics can change when computed per batch. See the
[data-adapter guide](mkdocs/guide/data-adapters.md) for whole-dataset blocks,
batching examples, and memory and failure behavior.

## Capture metadata and artifacts

Blocks receive the same stores as their pipeline. Use metadata for measurements
or context and artifacts for files or other outputs:

```python
from lgopy.core import Block, LgoPipeline, InMemoryArtifactStore, InMemoryMetadataStore

class SaveSummary(Block):
    """Record a dataset summary while preserving its input."""

    def call(self, values: tuple[float, ...]) -> tuple[float, ...]:
        """Save the count and input values.

        Args:
            values: Complete collection of measurements to record.

        Returns:
            The unchanged input tuple.
        """
        self.metadata.set("summary.count", len(values), attributes={"unit": "items"})
        self.artifacts.save(
            "summary/values.txt",
            "\n".join(map(str, values)),
            attributes={"artifact_type": "measurement_summary"},
        )
        return values

metadata = InMemoryMetadataStore()
artifacts = InMemoryArtifactStore()
pipeline = LgoPipeline.from_steps(SaveSummary(), metadata=metadata, artifacts=artifacts)
assert pipeline((1.0, 2.0)) == (1.0, 2.0)
assert metadata["summary.count"] == 2
assert artifacts.get_attributes("summary/values.txt") == {
    "artifact_type": "measurement_summary"
}
```

Use `FSSpecArtifactStore("file:///path/to/artifacts")` for filesystem output or an
appropriate remote URL for cloud storage. Store interfaces can also be implemented
by a host application. Custom attributes must be JSON-compatible.

Use unique keys when outputs from different steps or items must coexist.
Batching does not automatically clear in-memory stores, roll back database writes,
or delete files after a failure. See [stores](mkdocs/guide/stores.md) and
[store attributes](mkdocs/guide/store-attributes.md) for the persistence contract.

## Package and publish blocks

Build a folder or ZIP directly from a block defined in a Python source file:

```python
from normalize_block import Normalize

package = Normalize.build("build/normalize/1.0.0", format="zip", run_tools=False)
print(package.package_dir)
print(package.archive_path)
```

A package contains `block.py`, `__init__.py`, `requirements.txt`, `schema.json`,
and `manifest.json`. The build result includes inferred dependency pins and
static security findings. Enable optional formatter/linter tooling with
`run_tools=True`; `fail_on_security=True` rejects high-severity static findings.
These checks do not sandbox execution.

Publish the built directory to a local or remote fsspec-backed catalog, then
compose pipelines from stored packages:

```python
from pathlib import Path
from lgopy.core import LgoPipeline
from lgopy_catalog import BlockCatalog, FSSpecBlockStore

catalog = BlockCatalog(
    block_store=FSSpecBlockStore(Path("catalog").resolve().as_uri())
)
catalog.publish_package("build/normalize/1.0.0")
print(catalog.list_versions("normalize"))
print(catalog.search("normalize", category="numeric"))

pipeline = LgoPipeline.from_steps(
    catalog.create("normalize", version="1.0.0", scale=10.0)
)
assert pipeline([10.0, 20.0]) == [1.0, 2.0]
pipeline.save("catalog-pipeline.json")
restored = LgoPipeline.from_file("catalog-pipeline.json", catalog=catalog)
assert restored([30.0]) == [3.0]
```

Consumers need the catalog connection and block dependencies, not the publisher's
source repository. The catalog does not install dependencies automatically:
materialize selected packages and install their generated `requirements.txt` in
the execution environment. `load()` returns a class; `create()` returns a
configured instance. Use `remove_package(name, version)` to remove a package.

A catalog belongs to you or your organization; it is not a global public registry.
See [building packages](mkdocs/guide/building-blocks.md) and
[pipelines from published blocks](mkdocs/guide/catalog-pipelines.md).

## Discover blocks by meaning

`catalog.search()` provides text matching and metadata filters.
`catalog.semantic_search()` ranks indexed blocks against a natural-language query:

```python
# Requires a catalog configured with embeddings and a vector index.
matches = catalog.semantic_search("normalize numeric measurements", k=5)
for match in matches:
    print(match["name"], match["version"], match["distance"])
```

Built-in adapters support **Gemini**, **Ollama**, and **PostgreSQL/pgvector**;
custom embedding and vector-index implementations can be supplied through their
protocols. Indexing embeds manifest metadata, schemas, and the block's `call`
signature, docstring, and source.

Configure the embedding service and database, then publish or reindex packages
through that catalog. Adding embeddings does not automatically index previously
stored packages. Results include manifests and schemas; lower cosine distance
means a closer match. Inspect candidates and validate their configuration before
execution. See [semantic search](mkdocs/guide/semantic-search.md) and
[catalog adapters](mkdocs/guide/catalog-adapters.md) for complete setup examples.

### Experimental prompt-based assembly

`LgoPipeline.from_prompt(prompt, hub=..., llm=...)` can select blocks from a
runtime registry. Without an LLM it uses registry search and keyword matching;
a supplied compatible language model requires the corresponding LangChain and
provider dependencies. This is separate from catalog semantic search. Review the
selected order, default constructor arguments, and types before execution.

## Extend execution and file IO

**Callbacks:** pass callback objects through
`LgoPipeline.from_steps(..., callbacks=[callback])` to observe `on_start()`,
`on_step_start(step, X)`, `on_step(step, X)`, `on_end()`, and `on_error(error)`.
A plain object implementing these methods can report progress or collect timing
information. Use `lgopy.setup_logging(level="INFO")` to configure console logging.
Avoid the currently unavailable streaming callback import described below.

**Readers and exporters:** subclass `DataReader` or `DataExporter` and register
implementations with `DataReaderFactory.register(".extension")` or
`DataExporterFactory.register(".extension")`. Their `build(extension, **kwargs)`
methods create the selected implementation with your options. This provides a
common entry point for application-specific formats without adding format logic
to every processing block. See the [API reference](mkdocs/api-reference.md).

## Documentation

| Start here | Continue with |
| --- | --- |
| [Quickstart](mkdocs/quickstart.md) | [Blocks and pipelines](mkdocs/guide/blocks-pipelines.md) |
| [Data adapters and scope](mkdocs/guide/data-adapters.md) | [Metadata and artifacts](mkdocs/guide/stores.md) |
| [Block packaging](mkdocs/guide/building-blocks.md) | [Published-block pipelines](mkdocs/guide/catalog-pipelines.md) |
| [Semantic search](mkdocs/guide/semantic-search.md) | [Embedding and index adapters](mkdocs/guide/catalog-adapters.md) |
| [Runnable examples](mkdocs/examples/index.md) | [API reference](mkdocs/api-reference.md) |

Editable documentation lives in `mkdocs/`; generated GitHub Pages output lives
in `docs/`.

## Development and contributing

From the repository root:

```bash
uv sync --extra dev --extra docs
make check  # Syntax, typed-docstring coverage, pytest, and strict docs build
make build  # Build both workspace packages without publishing
```

Both packages share one release version and Git tag. Use
`python scripts/sync_versions.py 2.0.1` followed by `uv lock` to bump them together;
`make versions` checks their versions and peer requirements for drift.
Pushing a matching `v*` tag triggers the [release workflow](RELEASING.md), which
publishes both packages to PyPI and creates one GitHub release. Configure the
Trusted Publishers described in that guide before the first release.

Run individual checks with `make test`, `make docstrings`, or `make docs`.
See [CONTRIBUTING.md](CONTRIBUTING.md) for the workflow and
[SECURITY.md](SECURITY.md) for private reporting. Keep credentials and local
runtime data outside Git; [.env.example](.env.example) lists configuration keys.


## License

LgoPy and lgopy-catalog are licensed under [Apache-2.0](LICENSE).
