Metadata-Version: 2.5
Name: segmentstream-pipeline
Version: 0.1.0a56
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: anyio<5,>=4
Requires-Dist: certifi>=2024
Requires-Dist: dagster==1.13.16
Requires-Dist: httpx<1,>=0.28
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

Declare source tables with `s.ingestion` and warehouse SQL transformations with
`s.model`. The SDK resolves inputs, validates explicit partitioning and capabilities,
and compiles native Dagster definitions using the existing warehouse loader.

```python
from collections.abc import Iterator

import ibis.expr.types as ir
import segmentstream as s


@s.ingestion(
    tier=s.Tier.Bronze,
    partitioning=s.Partitioning.Unpartitioned,
    schema={"amount": "float64"},
)
def orders() -> Iterator[s.IngestionRow]:
    yield {"amount": 25.0}


@s.model(
    tier=s.Tier.Gold,
    partitioning=s.Partitioning.Unpartitioned,
    schema={"total_revenue": "float64"},
)
def revenue(orders: ir.Table) -> ir.Table:
    return orders.aggregate(total_revenue=orders.amount.sum())


defs = s.definitions(
    assets=[orders, revenue],
    jobs=[s.job("revenue", assets=[revenue], include_upstream=True)],
)
```

The [workspace documentation](../../docs/README.md) is organized by concept:
[connections](../../docs/connections.md), [ingestion](../../docs/ingestion.md),
[models](../../docs/models.md), [partitioning](../../docs/partitioning.md),
[jobs](../../docs/jobs.md), [execution](../../docs/execution.md), [checks](../../docs/checks.md),
[runtime values](../../docs/runtime.md), and [warehouse storage](../../docs/warehouse.md).
The CLI bundles these same pages into every starter workspace.

`s.validate(defs)` checks the graph and named job configuration without executing
assets. `s.describe(defs)` explains resolved inputs, per-asset policies, storage,
capabilities, same-day input mappings, and job selections without resolving
credentials. See the
[modular commerce starter](../../templates/workspace-app/pipeline/README.md)
for daily row ingestion, batched models, and a partitioned check.

Daily inputs may begin earlier or end later than their consumers when they use
the same timezone and day boundaries. The upstream must cover every downstream
date. The SDK supplies strict same-day mappings for inputs and ordering-only
dependencies; no native Dagster decorator or mapping option is needed. Named
jobs still select one daily calendar, with other calendars left as inputs.

## Identity-resolution support (0.1.0a39)

Models can declare `s.AllPartitions(asset)` inputs for retained-history rollups.
`WarehouseResource.materialize_scratch` and `iterate` execute bounded fixed-point
algorithms in warehouse SQL and retain their results until output publication.
The SDK owns scratch cleanup; each algorithm owns its convergence and key invariants.

## Compatibility

The SDK contains authoring, graph inspection, Ibis resources, connection HTTP
clients, and managed secret access. Managed orchestration, workers, PostgreSQL setup, and the metadata
webserver live in the internal [pipeline runtime](../pipeline-runtime/README.md).

In `0.1.0a28`, use `s.get_connection(key).get_http_client()` for provider requests.
The client forwards requests through the backend, where provider credentials and
refresh remain. The access-token and full-credential accessors have been removed.
See [connections](../../docs/connections.md) for declaring API destinations.

Workspace pipelines declare `segmentstream-pipeline[bigquery]`; our image supplies
the runtime. The old SDK `cloud-run` and `webserver` extras and
`segmentstream.runtime` module are removed in a22.

`segmentstream.dagster` retains its native exports. Native assets and high-level
declarations can coexist in `s.definitions`. Root `s.resource` and
`s.IngestionResource` remain available for legacy ingestion. Import their helpers
from `segmentstream`; the old public `segmentstream.ingestion` module is retired
so `s.ingestion` can be the decorator. [Migration examples](../../docs/native-dagster.md)
cover the mapping and IO defaults.

## Managed Google clients

The Sandbox runtime integration supplies refreshable credentials to the built-in
BigQuery warehouse and GCS IO manager. Ibis remains the user-facing warehouse API;
its internal official Google client renews credentials automatically throughout
the authorized run. Ten minutes is the lifetime of each access token, not a limit
on pipeline duration. Managed execution does not provide general ADC support.
See [runtime lifecycle and credential brokering](../../docs/sandbox-runtime.md).
Managed Dagster execution uses the backend compute provider API for every step,
including SQL assets and checks. The coordinator owns scheduling and retries; the
backend owns execution identity, scoped credentials and cleanup. See the
[compute provider contract](../../docs/provider-compute.md).

## Releases

SDK `0.1.0a39` adds explicit all-partition model inputs and managed warehouse
iteration with temporary SQL tables, supporting full-history identity resolution.

SDK `0.1.0a38` adds compatible daily calendars with strict same-day dependency
mappings. Consumers can start later than their upstream assets; each job still
selects one daily calendar.

The GA4 query client targets `0.1.0a37`. A model constructs it with
`s.ga4_bigquery_client(key, context=context)` and submits SQL through the backend;
BigQuery writes the result directly to the model's warehouse partition. See
[GA4 BigQuery connections](../../docs/bigquery-connections.md).
Publish the SDK before releasing a CLI whose starter
pins it. PyPI versions are immutable; do not republish an existing version.

The protected `pipeline-sdk-release.yml` workflow builds a wheel and source
distribution, then uses PyPI Trusted Publishing from the `pypi` GitHub environment. No
long-lived publishing token is stored in GitHub. Merging a new `version` in
`pyproject.toml` to `main` releases it and then creates its `pipeline-sdk-v<version>` tag; a
version already tagged or on PyPI is skipped. Pushing that tag, or running the workflow
manually, also releases the version. A pull request that changes `src/` or the published
`pyproject.toml` metadata must bump the version (`scripts/check-sdk-versions.mjs`); the app
and ontology SDKs follow the same rule and release the same way.

## Unreleased: full-refresh storage layout

Models with `s.Partitioning.Unpartitioned` may set
`storage=s.TableStorage(partition_by_date="session_date")`. They run as global
steps and atomically replace the whole table while retaining physical BigQuery
DATE partitions. Daily models keep selected-date replacement. `s.describe`
reports the storage layout independently of execution partitioning. See
[the model contract](../../docs/models.md#full-refresh-tables-with-physical-date-partitions).
Available in SDK `0.1.0a40`; update the workspace dependency and lockfile before using it.

## Ontology entity sources (0.1.0a56)

Pipeline manifest v10 describes assets, schemas, jobs and checks. Declare business measures in
[ontology entity sources](../../docs/ontology-aggregates.md). The semantic_models argument and
segmentstream.semantics package were removed. Release platform support before upgrading workspaces.

## Evaluation connection (0.1.0a43)

`EvaluationConnection(key="evaluation")` evaluates state against typed `boolean`,
`choice` and `score` questions through Vercel AI Gateway's evaluation API; SegmentStream
chooses the model and holds the credential, so it needs no configuration or
authorization. Call `s.evaluation_client("evaluation").evaluate(state=..., questions=...)`
in managed runs. See [evaluation connections](../../docs/evaluation-connections.md).
Connections inspection publishes connections manifest schema 8: upgrade the backend
before deploying declarations from this SDK.

## Own OAuth applications and LinkedIn Community (0.1.0a55)

`LinkedInAdsConnection`, `HubSpotConnection` and `GoogleSheetsConnection` accept
`application="own"`: users enter their own OAuth client ID and secret on the connection and
authorize through it, with the managed connector's scopes and read-only operations.
`LinkedInCommunityConnection(key="linkedin")` reads a LinkedIn member's post and follower
analytics and the pages they administer through the Community Management API; `scopes`
narrows the requested permissions to those its application has. Connections inspection
publishes connections manifest schema 12: upgrade the backend before deploying declarations
from this SDK. See [connections](../../docs/connections.md#linkedin-community).

## Apollo connection (0.1.0a54)

`ApolloConnection(key="apollo")` declares the Apollo.io API at
`https://api.apollo.io/api/v1/` with an API key, which the backend sends on the `x-api-key`
header. It is an `ApiKeyConnection` preset and publishes the existing connections manifest
schema 11. See [connections](../../docs/connections.md#apollo).

## Supabase connection (0.1.0a53)

`SupabaseConnection(key="supabase", project="<ref>")` declares one Supabase project's REST
API at `https://<ref>.supabase.co/rest/v1/` with a secret key (`sb_secret_`), which the
backend sends on the `apikey` header. It is an `ApiKeyConnection` preset and publishes the
existing connections manifest schema 11. See [connections](../../docs/connections.md#supabase).

## Stripe connection (0.1.0a52)

`StripeConnection(key="stripe")` declares the Stripe API with an API key, preferably a
restricted read-only key; the backend sends it as `Authorization: Bearer <key>`. It is an
`ApiKeyConnection` preset and publishes the existing connections manifest schema 11. See
[connections](../../docs/connections.md#stripe).

## Client credentials and Xero connection (0.1.0a51)

`ClientCredentialsConnection` declares OAuth 2.0 client-credentials access: users configure a
client ID and secret, and the backend exchanges them for bearer tokens. `XeroConnection`
reads one Xero organisation through a Xero custom connection; declare one key per
organisation. Connections inspection publishes connections manifest schema 11: upgrade the
backend before deploying declarations from this SDK. See
[connections](../../docs/connections.md#xero).

## Bearer values and Resend connection (0.1.0a50)

`ApiKeyConnection` headers accept `Bearer(Secret(...))`, so users configure only a token
and the backend sends `Bearer <token>`. `ResendConnection(key="resend")` declares the
Resend API with an API key. Connections inspection publishes connections manifest
schema 10: upgrade the backend before deploying declarations from this SDK. See
[connections](../../docs/connections.md#resend).

## Apify connection (0.1.0a49)

`ApifyConnection(key="apify")` declares access to the Apify API with a personal API
token. Deploy, paste the token into the connection, and use
`s.get_connection("apify").get_http_client()` with paths relative to
`https://api.apify.com/v2/`. It is an `ApiKeyConnection` preset, so it publishes the
existing connections manifest schema 9. See
[connections](../../docs/connections.md#apify).

## Job schedules (0.1.0a48)

`s.schedule(name, job=, cron=, timezone=, partitions=s.LatestPartitions(count))`
runs a named job on a cron schedule, at most once an hour. Register schedules with
`s.definitions(schedules=...)`. SegmentStream's scheduler, not Dagster, starts the runs.
Daily jobs run the newest `count` complete partitions on each tick. Inspection publishes
manifest version 9. Native Dagster `ScheduleDefinition`s now fail validation; they never
ran in SegmentStream. Upgrade the backend and the internal runtime before deploying
definitions from this SDK.

## Google Sheets connection (0.1.0a47)

`GoogleSheetsConnection(key="google_sheets")` declares read-only access to the
spreadsheets one Google account can open, through SegmentStream's Google app. Deploy,
authorize, and use `s.get_connection("google_sheets").get_http_client()` to read
spreadsheet metadata and values by spreadsheet ID. See
[connections](../../docs/connections.md#google-sheets). Upgrade the backend before
deploying declarations from this SDK.

## Local TLS roots (0.1.0a46)

Connection HTTP clients trust the Python installation's CA roots, as managed runtimes
require. When that store is empty, as in python.org macOS builds where
"Install Certificates.command" was not run, they also trust certifi's bundle, so
`segmentstream exec` works without setting `SSL_CERT_FILE`. A certificate that still
cannot be verified raises `RuntimeContextError` naming the verification failure, with the
original error as its cause; earlier SDKs report it as temporary unavailability.

## API-key connections (0.1.0a45)

`ApiKeyConnection(key=..., name=..., provider=..., headers={...}, query={...}, apis={...})`
declares HTTP access authenticated by static values such as API keys or tokens. Each
header or query parameter maps to a `Secret(...)` configured after deployment, or to
`secret_ref(...)`. The backend attaches the values to every proxied request; the
connection needs no authorization. Use `s.get_connection(key).get_http_client()` as
with OAuth connections. See [connections](../../docs/connections.md#declare-api-key-access).
Connections inspection publishes connections manifest schema 9: upgrade the backend
before deploying declarations from this SDK.

## Evaluation refusals (0.1.0a44)

When the evaluation service refuses requests outright, `evaluate` raises
`RuntimeContextRefusedError`, a Dagster `Failure` that disallows retries: the step fails
at once with the request reference instead of retrying. Earlier SDKs report it as
temporary unavailability. See [evaluation errors](../../docs/evaluation-connections.md#errors).

## HubSpot connection (0.1.0a42)

`HubSpotConnection(key="hubspot")` declares read-only access to one HubSpot
account's CRM through SegmentStream's HubSpot app. Deploy, authorize, and use
`s.get_connection("hubspot").get_http_client()`; the default `objects` API reads
contacts, companies and deals. See [connections](../../docs/connections.md#hubspot).
Upgrade the backend before deploying declarations from this SDK.
