Metadata-Version: 2.5
Name: segmentstream-pipeline
Version: 0.1.0a5
Summary: Dagster and Ibis runtime integration for SegmentStream pipelines
Project-URL: Homepage, https://segmentstream.com
Author: SegmentStream
License-Expression: Apache-2.0
License-File: LICENSE
License-File: NOTICE
Classifier: Development Status :: 3 - Alpha
Classifier: Intended Audience :: Developers
Classifier: Programming Language :: Python :: 3
Classifier: Programming Language :: Python :: 3.12
Classifier: Programming Language :: Python :: 3.13
Classifier: Typing :: Typed
Requires-Python: <3.14,>=3.12
Requires-Dist: dagster-postgres==0.29.16
Requires-Dist: dagster==1.13.16
Requires-Dist: ibis-framework<13,>=12
Requires-Dist: pydantic<3,>=2.10
Provides-Extra: bigquery
Requires-Dist: ibis-framework[bigquery]<13,>=12; extra == 'bigquery'
Description-Content-Type: text/markdown

# SegmentStream pipeline SDK

`segmentstream-pipeline` is the small runtime contract between a workspace's
Dagster definitions and the warehouse provisioned by SegmentStream. Dagster
continues to own assets and jobs, while Ibis continues to own relational
expressions. This package supplies lazy runtime configuration and durable
warehouse I/O.

The initial connector supports BigQuery, automatic dataset creation, full-table
replacement for unpartitioned assets, and native daily `DATE` partitioning. It
reads the following non-secret configuration when a pipeline first accesses the
warehouse:

- `SEGMENTSTREAM_WAREHOUSE_ENGINE`
- `SEGMENTSTREAM_WAREHOUSE_CATALOG`
- `SEGMENTSTREAM_WAREHOUSE_DEFAULT_NAMESPACE`
- `SEGMENTSTREAM_WAREHOUSE_LOCATION` (optional)

Configuration and authentication are deliberately lazy. Importing and
validating `definitions.py` during a deployment build does not connect to a
warehouse. In Cloud Run, the BigQuery connector uses the attached workload
identity through Application Default Credentials.

```python
import ibis
import ibis.expr.types as ir
import segmentstream.dagster as dg

from segmentstream import WAREHOUSE_IO_MANAGER_KEY, warehouse_resources


BRONZE_ORDERS = dg.AssetKey(["bronze", "orders"])
SILVER_ORDERS = dg.AssetKey(["silver", "orders"])


@dg.asset(
    key=BRONZE_ORDERS,
    io_manager_key=WAREHOUSE_IO_MANAGER_KEY,
    kinds={"ibis"},
)
def orders() -> ir.Table:
    return ibis.memtable(
        [{"order_id": "o-1", "amount": 100.0}],
        schema={"order_id": "string", "amount": "float64"},
    )


@dg.asset(
    key=SILVER_ORDERS,
    ins={"orders": dg.AssetIn(key=BRONZE_ORDERS)},
    io_manager_key=WAREHOUSE_IO_MANAGER_KEY,
    kinds={"ibis"},
)
def normalized_orders(orders: ir.Table) -> ir.Table:
    return orders.filter(orders.amount > 0)


defs = dg.Definitions(
    assets=[orders, normalized_orders],
    resources=warehouse_resources(),
)
```

Daily assets use Dagster's native daily partitions and declare the physical
BigQuery `DATE` column through SegmentStream metadata:

```python
from datetime import date

import ibis
import ibis.expr.types as ir
import segmentstream.dagster as dg

from segmentstream import (
    WAREHOUSE_IO_MANAGER_KEY,
    warehouse_asset_metadata,
)


daily = dg.DailyPartitionsDefinition(start_date="2026-01-01")


@dg.asset(
    key=["silver", "daily_orders"],
    partitions_def=daily,
    backfill_policy=dg.BackfillPolicy.multi_run(max_partitions_per_run=10),
    metadata=warehouse_asset_metadata(partition_by_date="event_date"),
    io_manager_key=WAREHOUSE_IO_MANAGER_KEY,
)
def daily_orders(context: dg.AssetExecutionContext) -> ir.Table:
    partition_dates = [date.fromisoformat(key) for key in context.partition_keys]
    return ibis.memtable(
        [{"event_date": value, "order_count": 0} for value in partition_dates],
        schema={"event_date": "date", "order_count": "int64"},
    )
```

SegmentStream accepts only unpartitioned assets and default-midnight
`DailyPartitionsDefinition` assets with `YYYY-MM-DD` keys. Deployment inspection
rejects other partition definitions and daily assets without physical partition
metadata. The IO manager maps Dagster's partition time window to a half-open
warehouse date range, filters upstream Ibis relations to that range, creates the
table with native daily partitioning on first materialization, and atomically
replaces only those dates on subsequent materializations.

Asset keys map to relations using a small convention:

- `["orders"]` uses the configured default namespace.
- `["bronze", "orders"]` uses the explicit `bronze` dataset.
- Other key shapes are rejected.

The workspace project is always supplied by SegmentStream and cannot be
overridden by an asset. Before writing an asset, the IO manager creates its
validated dataset with `CREATE SCHEMA IF NOT EXISTS` in the configured location.
This lets pipeline authors organize one workspace project into datasets such as
`bronze`, `silver`, and `gold` without provisioning them separately.

For local package development, install this project in editable mode rather
than adding a relative path dependency to a deployable pipeline.

Workspace pipelines declare only the SegmentStream SDK. It installs the pinned
Dagster and Ibis versions that belong to that SDK release:

```toml
[project]
dependencies = [
  "segmentstream-pipeline[bigquery]==0.1.0a5",
]
```

Pipeline definitions import `segmentstream.dagster` as their curated Dagster
namespace. Its objects are direct re-exports from Dagster, not wrappers. APIs
outside that namespace are not part of the SegmentStream Cloud compatibility
contract even if they remain importable from the underlying dependency.

The same installed package contains SegmentStream's private Cloud Run runtime:
the deployment inspector, persistent Dagster instance setup, and ephemeral
backfill coordinator. Workspace code does not call these modules directly.
Keeping them in this distribution ensures that the SDK, Dagster, and
`dagster-postgres` versions always move together; the backend only builds and
launches the installed runtime.

## Releases

Releases use the version declared in `pyproject.toml` and are published from the
`pipeline-sdk-v<version>` Git tag by the protected `pipeline-sdk-release.yml`
workflow. The workflow builds the wheel and source distribution in a job without
publishing credentials, then uses PyPI Trusted Publishing from the `pypi` GitHub
environment. No long-lived PyPI token is stored in GitHub.

PyPI releases are immutable. Increment the package version before creating a new
release tag; do not reuse a version that has already been uploaded.
