This release is a pre-release and may not be stable for production use.
SegmentStream pipeline SDK
segmentstream-pipeline contains the Python authoring surface and runtime
contract for a SegmentStream workspace. Dagster continues to own assets and
jobs, while Ibis continues to own relational expressions. This package supplies
declarative connection declarations, lazy runtime configuration, and durable
warehouse I/O.
Connections
Connections are workspace-level declarations and live in a separate top-level
connections/ component alongside app/ and pipeline/. Each custom OAuth
connection is executable by the isolated connection runtime, while importing
the declaration performs no authorization, network access, or secret resolution.
# connections/connections.py
from segmentstream.connections import OAuth2Connection, secret_ref
connections = (
OAuth2Connection(
key="google_ads",
name="Google Ads",
provider="google",
authorization_url="https://accounts.google.com/o/oauth2/v2/auth",
token_url="https://oauth2.googleapis.com/token",
client_id=secret_ref("google-oauth-client-id"),
client_secret=secret_ref("google-oauth-client-secret"),
scopes=["https://www.googleapis.com/auth/adwords"],
authorization_parameters={
"access_type": "offline",
"prompt": "consent",
},
),
)
The connection key is the stable identifier used by pipelines, the CLI, and
agents; name is its user-facing label.
Provider endpoints live in source. Both the client ID and client secret are
symbolic references supplied through the control plane and resolved only inside
the application's private connection runtime. Callback URLs, authorization
codes, refresh tokens, and access tokens are runtime values and never belong in
connection declarations. The runtime implements
authorization URL construction, authorization-code exchange, and refresh-token
exchange; the control plane owns OAuth state and durable encrypted token storage.
The separately published segmentstream-connections-runtime distribution runs
from the workspace's connections/ directory:
python -m segmentstream_connections_runtime inspect manifest.json
python -m segmentstream_connections_runtime serve
The service listens on PORT (default 8080) and exposes private internal
authorize, code-exchange, and refresh operations. References are injected as
namespaced environment values—for example, google-oauth-client-id maps to
SEGMENTSTREAM_CONNECTION_SECRET_GOOGLE_OAUTH_CLIENT_ID. Production Cloud Run
services must require IAM authentication; these token-bearing endpoints are not
a public workspace API.
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_ENGINESEGMENTSTREAM_WAREHOUSE_CATALOGSEGMENTSTREAM_WAREHOUSE_DEFAULT_NAMESPACESEGMENTSTREAM_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.
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:
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 explicitbronzedataset.- 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:
[project]
dependencies = [
"segmentstream-pipeline[bigquery]==0.1.0a8",
]
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.
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file segmentstream_pipeline-0.1.0a8.tar.gz.
File metadata
- Download URL: segmentstream_pipeline-0.1.0a8.tar.gz
- Upload date:
- Size: 98.4 kB
- Tags: Source
- Uploaded using Trusted Publishing? Yes
- Uploaded via:
twine/7.0.0 CPython/3.13.14
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
669efbdb9f2956155812fd042505451ea9b6aa96faddeeb6dcc4791e515588c6
|
|
| MD5 |
51f448a089bafaf18a109a9fa97acfde
|
|
| BLAKE2b-256 |
3d581b4cfeb79aac842425abe15288560b6f0bf30faf8add3520072e65c6fa81
|
Provenance
The following attestation bundles were made for segmentstream_pipeline-0.1.0a8.tar.gz:
Publisher:
pipeline-sdk-release.yml on segmentstream/segmentstream
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
segmentstream_pipeline-0.1.0a8.tar.gz -
Subject digest:
669efbdb9f2956155812fd042505451ea9b6aa96faddeeb6dcc4791e515588c6 - Sigstore transparency entry: 2684306678
- Sigstore integration time:
-
Permalink:
segmentstream/segmentstream@b55f55f34eeda9df57d9433fa5934a97f7bce2d8 -
Branch / Tag:
refs/tags/pipeline-sdk-v0.1.0a8 - Owner: https://github.com/segmentstream
-
Access:
private
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
pipeline-sdk-release.yml@b55f55f34eeda9df57d9433fa5934a97f7bce2d8 -
Trigger Event:
push
-
Statement type:
File details
Details for the file segmentstream_pipeline-0.1.0a8-py3-none-any.whl.
File metadata
- Download URL: segmentstream_pipeline-0.1.0a8-py3-none-any.whl
- Upload date:
- Size: 30.7 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? Yes
- Uploaded via:
twine/7.0.0 CPython/3.13.14
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
7a99d12ae32157909b7b664f1d35bc114221fd6d279f2118717e31603213f121
|
|
| MD5 |
ec8d1c8e21e0aab2fcf7437e190181ea
|
|
| BLAKE2b-256 |
4b609602a0e966791cd210d860b3f84b36e8ebc26ca26cad6666915f8ee77196
|
Provenance
The following attestation bundles were made for segmentstream_pipeline-0.1.0a8-py3-none-any.whl:
Publisher:
pipeline-sdk-release.yml on segmentstream/segmentstream
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
segmentstream_pipeline-0.1.0a8-py3-none-any.whl -
Subject digest:
7a99d12ae32157909b7b664f1d35bc114221fd6d279f2118717e31603213f121 - Sigstore transparency entry: 2684306701
- Sigstore integration time:
-
Permalink:
segmentstream/segmentstream@b55f55f34eeda9df57d9433fa5934a97f7bce2d8 -
Branch / Tag:
refs/tags/pipeline-sdk-v0.1.0a8 - Owner: https://github.com/segmentstream
-
Access:
private
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
pipeline-sdk-release.yml@b55f55f34eeda9df57d9433fa5934a97f7bce2d8 -
Trigger Event:
push
-
Statement type: