Skip to main content

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 · Quickstart · Features · Documentation · 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:

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:

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 for all extras and their current limitations.

Quickstart

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

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:

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.

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:

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 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:

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 and store attributes for the persistence contract.

Package and publish blocks

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

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:

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 and pipelines from published blocks.

Discover blocks by meaning

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

# 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 and catalog adapters 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.

Documentation

Start here Continue with
Quickstart Blocks and pipelines
Data adapters and scope Metadata and artifacts
Block packaging Published-block pipelines
Semantic search Embedding and index adapters
Runnable examples API reference

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

Development and contributing

From the repository root:

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, 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 for the workflow and SECURITY.md for private reporting. Keep credentials and local runtime data outside Git; .env.example lists configuration keys.

License

LgoPy and lgopy-catalog are licensed under Apache-2.0.

Metadata

Release files for lgopy 2.0.0

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for lgopy 2.0.0
File Size Uploaded
lgopy-2.0.0.tar.gz 480.7 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for lgopy 2.0.0
File Interpreter ABI Platform
lgopy-2.0.0-py3-none-any.whl Python 3 none any Details

Total release size: 536.2 kB

Release files / lgopy-2.0.0.tar.gz

Download URL lgopy-2.0.0.tar.gz
Size 480.7 kB
Tags Source
SHA-256 checksum
How to use checksums
8203a2742193625a3a11203c9a3f576aa61687ff5c76a9463e11edcab1fc09d1
BLAKE2b-256 checksum
How to use checksums
edb9646de5442966656c0970fd6bbd94c614129d96795bd02d04df6be365f918
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/7.0.0 CPython/3.13.14

Provenance

Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.

PyPI Publish Attestation

PyPI verified that this artifact, at this checksum, originated from the publisher listed below.

Signed by GitHub Actions, verified by PyPI on Sep 14, 2026.

Transparency log

Release files / lgopy-2.0.0-py3-none-any.whl

Download URL lgopy-2.0.0-py3-none-any.whl
Size 55.5 kB
Tags Python 3
SHA-256 checksum
How to use checksums
b1c7d0fc21f0ce9369a965c904c05b203592bc1c7caba1b4fd580562b8480eef
BLAKE2b-256 checksum
How to use checksums
970be6dacd8cd44c4df991e18c79010d8aca617b465299b83fb003805deb117e
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/7.0.0 CPython/3.13.14

Provenance

Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.

PyPI Publish Attestation

PyPI verified that this artifact, at this checksum, originated from the publisher listed below.

Signed by GitHub Actions, verified by PyPI on Sep 14, 2026.

Transparency log

Release history Release notifications | RSS feed

This release

2.0.0 This release

2 release files

1.5.8

2 release files

1.5.7

2 release files

1.5.6

2 release files

1.5.4

2 release files

1.5.3

2 release files

1.5.2

2 release files

1.5.1

2 release files

1.5.0

2 release files

1.0.1

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page