Skip to content

Value Stream — High-Level Architecture

Document Architecture overview
Companion docs design/replacement-design.md (detailed), concepts/domain-model.md, reference/processors.md, reference/algorithms.md, reference/readers-and-formats.md, reference/expression-dsl.md, reference/chart-catalog.md, reference/faq.md
Audience Engineers, architects, and senior stakeholders
Current stack Polars · DuckDB · Streamlit · Plotly · Apache DataSketches · PyArrow Parquet
Current headless surfaces Read-only FastAPI HTTP API · local stdio MCP
Deferred surfaces Remote HTTP MCP · OIDC/multi-user service deployment

1. Mission and one-line description

Value Stream is a configuration-driven, aggregate-first business intelligence platform for marketing, ML, and customer-lifecycle metrics. It ingests batch exports from upstream operational systems (typically Pega CDH Interaction History and Product Holdings), reduces them to small, mergeable sufficient statistics, and serves business reports and dashboards from those persisted aggregates — never from source rows or processor checkpoints.

Everything that can be expressed as a small, fixed-size summary per group-by tuple is in scope. A processor may declare a finite dependency window when that summary requires sequence context and may retain minimal, rebuildable internal state to avoid repeated source scans. Unbounded event histories, report-visible identity sets, and arbitrary per-entity feature serving remain out of scope (or are approximated via sketches, or moved to a snapshot processor).

2. Design forces and quality attributes

Quality attribute What it means here How the architecture satisfies it
Aggregate-first serving Reports and queries must not process source or identity-level processor rows Each processor publishes sufficient statistics; every read surface uses the governed aggregate query layer
Bounded internal state An exact bounded processor may avoid repeated source scans without creating another business source Optional state is minimal, versioned, hash-sharded, transactionally reconciled, independently retained, and rebuildable from authoritative input
Configurable Non-developers can add metrics, dashboards, and group-by columns YAML DSL with JSON-Schema validation; closed expression AST replaces eval-strings
Deterministic Same computation contract + same file fingerprints = same numbers computation hash covers workspace defaults, source behavior, and processor semantics; merge ops are associative-commutative
Idempotent ingestion Unchanged chunks are skipped; changed files are reprocessed (source, chunk, source-computation-hash, file_hash) planning; the globally latest successful chunk attempt governs reads and reuse
Observability first Every number is traceable to a chunk and a config pipeline_run_id, chunk_id, period, created_at, config_hash columns on every aggregate row
Multi-grain queries Same metric available at its base grain and every coarser calendar grain Each processor stores one base-grain aggregate; the query layer merges it to the requested coarser grain
Operates on one node Targets ~10–100 GB raw input per workspace Polars + DuckDB; chunk-level concurrency; optional shard-by-source
Friendly defaults First-time users get a working dashboard out of the box Built-in processor presets, sample workspace, demo dataset, generated dashboards
Stack continuity Reuse Polars, DuckDB, Streamlit, Plotly Parquet remains the business-aggregate rest format; bounded rolling DuckDB state accelerates exact ingestion and FastAPI/local MCP reuse the governed query layer

Non-goals: streaming/CDC, distributed compute, ad-hoc warehouse SQL on raw events, replacing Pega CDH semantics.

3. Stakeholders and use cases

  • Marketing analyst opens Marketing Overview and asks "what was CTR yesterday on Web/Leaderboard for Cards/Loans?" The dashboard tile resolves to a daily engagement aggregate; the planner picks one physical Parquet partition; the answer comes back in milliseconds.
  • Data engineer schedules a daily ingestion run via CLI/cron and monitors freshness in the Ops page.
  • Decision scientist authors a new metric (Cost_per_Impression = Cost / Impressions) by editing metrics.yaml, validates with valuestream validate, and previews it in the Builder UI without touching Python.
  • Product owner opens the Experiments page and reads the chi-square p-value and odds-ratio CI for the latest A/B test.
  • LLM agent answers a question through Chat With Data or the local MCP server by calling aggregate metric tools such as metric_query(metric="CTR", group_by=["Channel"], ...). Raw data is never exposed.
  • Auditor asks "where did this number on the dashboard come from?" and gets a chunk list + config hash + YAML body in three clicks.

4. C4 — Context

                                +-------------------------------------+
                                |          Value Stream Workspace          |
                                |  (one logical environment, e.g.    |
                                |   "BDT", "RBB", "Demo")            |
                                +------------+------------------------+
                                             ^
                                             |
+--------------------+    files    +---------+----------+    config    +----------------+
| Upstream operational| ----------> | Value Stream Platform   | <-----------| Configuration  |
| systems (Pega CDH,  |             | (this design)      |             | YAML in git    |
| product holdings,   |             +---------+----------+             +----------------+
| subscriptions)      |                       ^
+--------------------+                        |
                                              | reads (UI / SDK / SQL export)
                              +---------------+----------------+
                              |                                |
                       +------+------+                  +------+------+
                       |  Streamlit  |                  |  Notebook,  |
                       |  UI         |                  |  SDK, SQL   |
                       +-------------+                  |  export     |
                                                        +------+------+
                                                               |
                                                 local MCP + read-only HTTP API;
                                                 remote HTTP MCP deferred

5. C4 — Containers (logical)

+-------------------------+   +-------------------------+   +-------------------------+
| 1. Configuration store  |   | 2. Aggregate store      |   | 3. Metadata store       |
| ----------------------- |   | ----------------------- |   | ----------------------- |
| YAML files in git or    |   | Parquet files,          |   | DuckDB databases:       |
| local catalog/ folder.  |   | hive-partitioned by     |   |   chunks, runs,         |
| JSON-Schema validated.  |   | period under            |   |   config_versions,      |
| Hashed -> config_hash.  |   | aggregates/<src>/<proc>/|   |   lineage.              |
+-----------+-------------+   |   <grain>/period=...    |   +-----------+-------------+
            |                 +-----------+-------------+               |
            |                             |                             |
            v                             v                             |
+-------------------------+   +-------------------------+                |
| 4. Ingestion engine     |-->| 5. Query layer          |<---------------+
| ----------------------- |   | ----------------------- |
| Discovery, grouping,    |   | Planner picks physical  |
| chunked Polars pipeline,|   | aggregate; executor     |
| processor fan-out,      |   | runs metric DSL formulas|
| compaction, ledger      |   | and sketch queries.     |
| writes.                 |   +-----------+-------------+
+-----+-------------------+               |
      ^                                   |
      |                                   v
      |                         +-------------------------+
      |                         | 6. Surfaces             |
      |                         | ----------------------- |
      |                         |  - Streamlit UI         |
      |                         |  - Python SDK           |
|                         |  - SQL via DuckDB views |
|                         |  - local MCP tools      |
|                         |  - read-only HTTP API   |
|                         |  - remote MCP deferred  |
      |                         +-------------------------+
      |
      |
+-----+-----+
| 7. CLI    |
| (valuestream) |
| run/     |
| migrate/ |
| vacuum/  |
| validate |
+----------+

Container responsibilities:

  1. Configuration store — YAML files versioned in git: pipelines.yaml, processors.yaml, metrics.yaml, and dashboards.yaml. Loaded at startup, validated, hashed.
  2. Aggregate store — file-system-rooted Parquet directories with hive partitioning. The only persisted data read by business query surfaces. A separate processor-state namespace may contain non-queryable, source-derived acceleration state for a bounded processor.
  3. Metadata store — small DuckDB databases tracking chunks, runs, config versions, lineage. Single-writer per database file; readers can be many.
  4. Ingestion engine — turns files into aggregates. Reads files lazily via Polars, applies transforms, fans out to processors, writes Parquet partials, runs compactions, updates the chunks ledger.
  5. Query layer — turns metric requests into reads from the aggregate store. Plans, scans Parquet via DuckDB, materializes Polars frames, applies derived metric DSL, returns rows.
  6. Surfaces — current read clients on top of the same query layer are Streamlit UI, Python SDK, DuckDB export, Chat With Data, local stdio MCP, and a read-only FastAPI HTTP API. Remote HTTP MCP and multi-user/OIDC deployment are deferred.
  7. CLI — operator entry point: run, validate, backfill, vacuum, serve.

The Streamlit reports surface reads typed page-filter definitions from dashboards.yaml. Omitted filters or [] mean no filters; no filters are inferred from processors. Because Value Stream is aggregate-first, only persisted processor group_by columns are eligible. Each filter declares all_tiles or compatible_tiles coverage. The capability matrix is validated, partial coverage is shown on the active chip, and every unsupported tile names the active filters it did not apply; filters are never skipped silently.

Before scanning Parquet, the query executor uses file lineage to select only partials whose config_hash matches the current Processor computation hash. It never blends files from different Processor schemas or substitutes nulls for newly configured aggregate states. If no current-hash partial has been published, the query reports the metric as not ready and the Reports UI shows Backfill required instead of exposing a storage-library exception. Legacy imports without file-lineage rows retain the embedded-hash scan fallback. Freshness reads use the same current-hash lineage selection before opening Parquet, so status badges cannot fail by asking Polars to reconcile superseded and current Processor schemas. The report summary labels this upstream result as the aggregate status; chart-rendering failures are reported on the affected tile and are not presented as page-health success.

Metric display metadata lives in metrics.yaml; page controls, KPI semantics, and tile presentation live in dashboards.yaml. Typed presentation properties cover labels, units, value formats, favorable direction, explicit KPI-strip placement, previous-period/target comparison, sparklines, absolute/index/change scales, semantic category colors, independent line-dash and marker-symbol dimensions, goal/reference lines, sorting, and conditional colors. These settings are consumed by the query/presentation boundary and do not alter persisted aggregate state or metric calculation semantics.

For external BI tools, valuestream export-duckdb <workspace> --grain <grain> creates a materialized DuckDB file with one table per metric at the selected grain. Each table is populated through the normal metric query layer using all persisted group_by dimensions from the metric's processor, so SQL consumers see ordinary metric-output columns rather than serialized sketch state. This is an export artifact, not the canonical store; Parquet aggregates remain the source of truth.

6. C4 — Components inside the ingestion engine

+--------------------------------------------------------------+
|                    Ingestion Engine                          |
|                                                              |
| +------------+   +-----------+   +-----------+   +---------+ |
| | Discovery  |-->| Reader    |-->| Transforms|-->|Processor| |
| | & grouping |   | (parquet, |   | pipeline  |   |fan-out  | |
| +------------+   |  pega zip,|   +-----------+   +----+----+ |
|                  |  csv,xlsx)|                        |      |
|                  +-----------+                        v      |
|                                              +--------+----+ |
|                                              | Per-processor| |
|                                              | chunk        | |
|                                              | aggregator   | |
|                                              +--------+----+ |
|                                                       |      |
|                                                       v      |
|                                              +--------+----+ |
|                                              | Partial      | |
|                                              | Parquet      | |
|                                              | writer       | |
|                                              +--------+----+ |
|                                                       |      |
|                                                       v      |
|                                              +--------+----+ |
|                                              | Compactor    | |
|                                              | (daily ->    | |
|                                              |  monthly ->  | |
|                                              |  summary)    | |
|                                              +--------+----+ |
|                                                       |      |
|                                                       v      |
|                                              +--------+----+ |
|                                              | Chunk ledger | |
|                                              | + run record | |
|                                              +--------------+ |
+--------------------------------------------------------------+

7. C4 — Components inside the query layer

+--------------------------------------------------------------+
|                       Query Layer                            |
|                                                              |
|   +-------------+   +---------+   +----------+   +--------+  |
|   | Resolver    |-->| Planner |-->| Executor |-->| Metric |  |
|   | (metric ->  |   | (pick   |   | (DuckDB  |   | DSL    |  |
|   | processor + |   | grain & |   | scan +   |   | (formula|  |
|   | states)     |   | filters)|   | Polars)  |   | / curve|  |
|   +-------------+   +---------+   +----------+   |  / test|  |
|                                                  |  / RFM)|  |
|                                                  +--------+  |
|                                                              |
|   Cross-cuts: cache (LRU keyed by                            |
|     metric_id, dim_set, filter_hash, grain, config_hash),    |
|     freshness reporter, lineage emitter.                     |
+--------------------------------------------------------------+

8. End-to-end data flow

For sources with materialize_transforms: true, the reader/transform portion below is collected once per chunk (with the Polars streaming engine when reader.streaming: true). Processor lazy plans then fan out together from that shared eager frame with the in-memory engine. The shared transformed frame is released after fan-out, and each processor/grain frame is released after its immutable aggregate is written. This is an execution strategy only: no raw or transformed row is persisted by ordinary processors, and streaming/materialization settings do not change the computation hash.

An explicitly bounded lookback processor is the one exception to current-only processor input. For frequency_response, the dependency planner selects the target plus the preceding daily chunks needed to cover its configured window and optional partition-lag allowance. The allowance widens only that closure; timestamp predicates retain the exact semantic window. The target chunk's ledger fingerprint covers the complete raw dependency closure, and ordinary processors still consume only the target chunk.

With checkpoint.mode: source_scan, the runner transforms that closure into an ephemeral current/history frame for the frequency processor. With checkpoint.mode: persistent_sharded, schema revision 8 opens the one stable rolling.duckdb addressed by source and processor and keeps its connection alive for the complete source run. Revision 8 is the default and only supported checkpoint schema. Targets execute in ascending ISO-date order. Polars streams the complete prepared current target through the Arrow C Stream interface into a temporary DuckDB relation; it is never persisted. The database retains only bounded exposed-rank-1 history with source chunk id and logical customer shard, plus a transactional journal of authoritative raw fingerprints.

For each logical shard, native DuckDB SQL combines current candidates with the rolling history and performs cross-day contact normalization, the strict frequency window, and selected rank-2 resolution. It returns only enriched target rank-1 rows to Polars through Arrow; Polars applies configured state predicates, count/value-sum aggregation, grouping, and provenance. Appending the new narrow history, journal fingerprint, and retention deletes is one state transaction. That acceleration commit may precede aggregate publication but never authorizes query visibility or chunk reuse. Before each pending target, the journal is reconciled with that target's expected ordered history closure. Rows outside that closure are removed. A missing suffix of the remaining exact fingerprinted prefix—including chunks skipped by aggregate-ledger reuse—is filled from IH without publishing aggregates. Fingerprint/order/non-prefix mismatch resets the required closure. Valid but incompatible metadata replaces the database at the same stable path and rebuilds it from authoritative IH; older schema revisions are not migrated. Corrupt, identity-tampered, or dtype- drifted state fails closed, and --force replaces it. Expired history is pruned after every chronological source-day commit. Every 30 commits, CHECKPOINT runs on the same open connection; the writer closes once at the source-run boundary. Both execution modes publish the same mergeable grouped states, and reports cannot address the state namespace. The DuckDB connection remains on the owning thread: one shard uses a batched lazy stream, while multi-shard execution overlaps a main-thread fetch with one worker's Polars aggregation over a detached Arrow table and keeps at most two complete shard results live. Both paths canonicalize decision time to timezone-naive UTC whole seconds before normalization; reporting-day derivation reattaches UTC before calendar-timezone conversion. This is specified by ADR 0007 and ADR 0008.

Order-sensitive bounded ML samples are the exception that proves the execution rule: when a score processor needs personalization or novelty, ingestion adds a temporary scan-order index before transforms and uses it inside those callbacks. That keeps results invariant across streaming and in-memory scheduling without persisting the index.

[Files in source folder]
      |
      v
[Discovery]  -- glob + group_by_filename pattern --> chunk_id
      |
      v
[Pipeline run row: status=running] -- durable publication barrier
      |
      v
[Reader]     -- pega_ds_export | parquet | csv | xlsx --> Polars LazyFrame
      |
      v
[Transforms] -- rename_capitalize, parse_datetime, derive_calendar,
                derive_action_id, filter (AST), dedup, defaults
      |
      v
[For each processor bound to this source:]
      |
      v
   [Processor.chunk_aggregate]
      |    -- group_by(group_by columns + time_grain)
      |    -- aggregate states (counts, sums, mins, maxes, sketches)
      v
   [Base chunk aggregate state]
      |
      v
[For each configured grain: processor.compact(base state)]
      |    -- remove finer calendar keys
      |    -- merge states by state-type rule, without rereading raw rows
      v
[Immutable run/chunk parquet partials written atomically]
      |    -- writer returns path/hash/rows/size/timestamp receipts
      |
      v
[Lineage transaction committed] -- receipts for every written aggregate path
      |
      v
[Chunk ledger row status=ok] -- last durable chunk commit marker
      |
      v
[Pipeline run finalized] -- ok / failed; recovered interruptions may be partial
      |
      v
=== ingestion done ===

=== query path begins ===

[Tile / SDK call: metric=CTR, grain=Day, fields=[Day,Channel,Group], time=2024-08]
      |
      v
[Resolver] -- CTR -> processor=engagement, states=[Positives,Negatives]
      |
      v
[Planner]  -- pick aggregates/ih/engagement/daily, period in [2024-08]
      |
      v
[Executor] -- DuckDB read_parquet(...) + compatible tile/page filters
      |
      v
[Polars frame]
      |
      v
[Metric DSL apply] -- CTR = Positives / (Positives + Negatives)
      |
      v
[Return rows + query provenance (stored grain, catalog/computation hashes,
 run IDs, chunk IDs, aggregate scan count, latest created_at)]

9. Storage layout

<workspace>/
├── catalog/                # versioned YAML configs
│   ├── pipelines.yaml
│   ├── processors.yaml
│   ├── metrics.yaml
│   └── dashboards.yaml
├── aggregates/             # the only place business data lives
│   └── <source_id>/
│       └── <processor_id>/
│           └── <grain>/
│               └── period=YYYY-MM/part-<run>-<chunk>.parquet
├── .valuestream/
│   └── state/              # non-queryable, rebuildable processor state
│       └── frequency_response/
│           └── source=<source_id>/
│               └── processor=<processor_id>/rolling.duckdb
├── meta/                   # metadata DBs (DuckDB)
│   ├── chunks.duckdb
│   ├── pipeline_runs.duckdb
│   ├── config_versions.duckdb
│   ├── lineage.duckdb
│   └── aggregate_views.duckdb  # governed views over successful aggregates

The processor-state namespace is deliberately outside aggregates/ and meta/: DuckDB views, metric planning, lineage publication, and chunk ok markers never expose it. A persistent frequency checkpoint has one stable path per source and processor; schema, hash, Polars, config, and layout values do not create directory levels. Schema and hashing revisions, Polars version, processor computation hash, shard count, history projection, and customer dtype stay inside the database as compatibility metadata; DuckDB version is retained for audit. The transaction journal records chunk ids and raw fingerprints.

The writer retains at most checkpoint.retention_days source-day journal entries and reconciles the journal before every pending target. Retention cannot be smaller than ceil((window_hours + partition_lag_hours) / 24); the default 168-hour window with zero lag retains only the last seven calendar source days, with fewer stored chunks when dates are missing. Expired rows are deleted after every chronological source-day commit. Every 30 commits, DuckDB CHECKPOINT folds the expected WAL and partially reclaims or reuses deleted-row space on the same open connection; it does not guarantee complete compaction or immediate file shrinkage. The connection closes once at the source-run boundary. Missing or valid-but-incompatible state is initialized from source IH at the same path; corrected fingerprint/order state is rebuilt deterministically. Corrupt or identity-invalid state fails closed and can be replaced with a forced rebuild.

The same workspace layout is used for every variant (BDT, RBB, NBS, Demo, …). Variants are separate workspaces; there is no commingling of variants inside one workspace.

10. Configuration model

Four YAML files form the catalog; each has a published JSON Schema in schemas/.

File What it defines
pipelines.yaml Sources (where files come from) + readers + transforms + defaults
processors.yaml Processors bound to sources, with group_by, one time.grain, states, outcome rules
metrics.yaml Derived metric definitions (formula, sketch query, variant compare, contingency test, …)
dashboards.yaml Dashboards, pages, tiles; tiles bind to metrics, not processors

The full DSL is specified in design/replacement-design.md §7 and concepts/domain-model.md §3. Expression semantics are in reference/expression-dsl.md.

Three catalog identities plus internal checkpoint compatibility metadata serve different boundaries:

  • the catalog hash identifies the complete authored catalog for audit, including execution/storage settings such as checkpoint.*;
  • the source computation hash covers workspace defaults, the source reader/schema/transforms/defaults, and the result-affecting semantics of all processors bound to that source, and controls chunk skip/reprocess decisions;
  • the processor computation hash covers workspace defaults, source behavior, and one processor's result-affecting semantics, and is persisted as aggregate config_hash; For frequency_response, checkpoint mode, shard count, and retention are excluded from processor and source computation hashes. Mode selects execution, shards select logical SQL partitioning inside the stable database, and retention selects bounded live-state policy. Changing them therefore leaves unchanged aggregate chunks skipped; incompatible state is rebuilt lazily at the same path for the next bounded target. Window, partition-lag, contact, alternative grouping, filter, published grouping, and state changes remain semantic and invalidate the appropriate aggregates.

Presentation-only descriptions, dashboards, and metric prose do not invalidate ingestion. Canonical payloads for the three catalog/computation identities are inserted into meta/config_versions.duckdb; each emitted aggregate file is recorded in meta/lineage.duckdb with its run, chunk, processor, grain, period, hash, row count, and path.

The Parquet writer derives that lineage receipt from the in-memory partition while writing it; before the chunk marker, the ledger confirms the file still exists with the recorded size. The normal ingestion path does not reopen the new file merely to recover metadata it already knows. Recovery still deep-scans embedded provenance before publishing an interrupted run.

Data Load dispatches source, workspace, and clean-rebuild actions to daemon threads in the Streamlit server process. A locked, process-local registry keyed by resolved workspace and run scope supplies poll-friendly progress across browser reloads and websocket reconnects. It is not a durable job queue: restarting the server ends those threads, after which the normal ingestion ledger recovery verifies and adopts completed chunk work.

Run preflight still validates structural catalog invariants and aggregate contracts. It does not reject an expression merely because its physical source field was absent from the catalog-inferred schema: the reader and transform executor validate those references against each real input chunk. Sample-backed authoring remains the earlier, friendlier place to catch missing physical fields.

Configuration authoring surfaces

Value Stream exposes a top-level Build choice over two Streamlit authoring paths backed by the same YAML catalog. Start from a sample enters AI Configuration Studio; Configure the current workspace enters Configuration Builder. The landing page chooses a workflow but does not create a third authoring store.

  • Configuration Builder is the catalog-first, validation-first editor for the active workspace. Its compact outline covers workspace health, sources, processors, dimensions, metrics, reports/tiles, chat review, settings, and Export current workspace. Each object editor compares a canonical session-local revision with the applied object. Internal navigation preserves a dirty revision until the user applies or discards it; simply visiting a step cannot make it dirty. A current object exposes exactly one Apply to workspace action, and apply never starts ingestion. Source steps create new sources or edit existing ones in place, including reader runtime settings, schema keys, defaults, dataset filters, calculated fields, and the in-memory materialize_transforms execution toggle; create mode rejects an existing Source ID rather than replacing it. Create mode also exposes an explicit, optional bounded sample inspection after reader settings are entered. Opening the editor performs no read, the preview persists only field names and types in session state, and the observed names seed later source-field controls without changing catalog configuration or starting ingestion. An OutcomeDate or OutcomeTime timestamp in create mode deterministically seeds the standard Pega datetime, calendar, action-ID, and numeric-cast transforms, constrained by observed optional fields when a sample was loaded. Filters are authored either as rule rows (field, operator, value, enabled) or editable, validated raw expression-AST YAML; calculated fields become derive_column transforms with typed AST expressions. Processor steps edit group-by dimensions, one physical base grain, states, and optional processor filters. Metric editors show the computed output contract as a read-only list. Report controls that consume metric-owned measures select only from that list and persist the choice as metric_output; dimension roles remain separate. Page-settings writes merge into the existing page and preserve its tiles, dashboard layout, and theme. Chat review shows which aggregate metrics will be available to Chat With Data and edits chat-only LLM prompt/description guidance in ai.yaml; settings edit workspace defaults plus dashboard theme.
  • Digest state authoring keeps storage and query roles separate without making the user repeat them. An unconditioned t-digest or KLL state is the persisted mergeable binary accumulator. Every structured processor-write path persists a quantile-less distribution metric in the same catalog transaction when no equivalent state binding exists. That metric supplies the stable public query/report identity and full distribution output; additional named percentiles remain optional. Outcome-specific positive and negative digests remain internal curve/calibration inputs and are not auto-exposed.
  • Configuration Builder checkpointing persists only privacy-filtered, JSON-safe draft registry state, the current step, a UTC timestamp, and the full base-catalog hash in meta/config_builder_checkpoint.json. Prompt, credential, provider, sample/upload, raw-provider, bytes, and DataFrame state is never written; a draft that would become incomplete is omitted. Restore remains subject to the object baseline and validation gates, catalog drift requires explicit reconciliation, and the file expires after seven days or is deleted when the recoverable registry becomes empty.
  • AI Configuration Studio is the sample-first, optionally model-assisted draft workflow. Uploaded bytes are preview-only; the user separately reviews the generated production source plan. CSV, Parquet, JSON, and explicitly detected Pega/archive samples receive format-specific reader defaults, while unsupported archive shapes fail before a misleading source can be drafted. Required-field mappings select only fields present in the approved schema. Sample values are excluded from model prompts by default. Approved schema names, types, null counts, and unique counts remain part of the prompt, while hidden field names are excluded. Every sample-backed model action requires confirmation of the exact sample, provider, model, endpoint route, approved fields, and example-value scope; changing that scope invalidates the confirmation and its widget state, while ordinary step navigation preserves it. The always-open governed Copilot is the first configuration surface after that confirmation and before the active step's manual controls. Effective post-transform field names are authoritative for source filters and every downstream calculation, processor, metric, and report; observed physical sample columns seed catalog validation and evolve through the declared transform order. A draft whose active-source naming transform disagrees with the current Sample contract is retained for comparison but cannot call a provider or Apply; deterministic regeneration presents the complete source naming change and dependent consumers as one dependency-closed review bundle. User-initiated model work preflights the configured provider/model/credential capability, reports a safe corrective action, and caches a successful check for the session. Generation parses, merges, and validates a candidate before review; a bounded sequence of at most two repair passes may run inside the same named operation. An unrecoverable candidate is discarded while the last valid draft remains available, with deterministic generation as the validated fallback. Review uses semantic, dependency-closed bundles in the main canvas; removals require explicit selection and exact YAML stays collapsed. Governed Copilot remains available for read-only explanation while bundles are pending but blocks mutations that could overwrite them. Applying a reviewed revision validates and writes sources, processors, metrics, dashboards, and ai.yaml inside one rollback boundary, preserving theme, layout, page, tile, and chat-guidance properties.

Both workflows display one revision ledger: editing draft → ready for review → reviewed → applied → data refresh required or report ready. Validation is keyed by the canonical revision, so a field change invalidates the prior verdict and review. Current-workspace and draft verdicts are never presented as if they describe the same object.

Every Builder catalog mutation and its post-write validation share one rollback boundary. A failed write or invalid resulting workspace restores all affected catalog files before control returns to the UI.

Both metric steps also expose one shared, versioned KPI recipe library. The packaged recipe YAML is an inert authoring artifact: it describes business questions, aggregate capability requirements, bounded install parameters, closed processor-state/metric templates, method accuracy, and report recommendations. A deterministic resolver reports ready, mapping_required, or backfill_required; only an explicit install materializes ordinary processor/metric/tile YAML. AI Studio installs into its draft, while Configuration Builder installs into the active catalog. See the KPI recipe reference.

Recipe mapping is business-facing: sketch-backed metrics select a processor-owned field and any algorithm declared compatible by the recipe; paired digests select one score field, and funnels select stages/populations. Recipes may also expose bounded values such as a relative materiality threshold. State IDs and aggregate parameters remain technical detail. A missing field/algorithm state or a parameter-resolved, recipe-authored filtered state becomes a deterministic processor-state proposal. Before installation, both Studios show the exact generated YAML patch plus the named source, fields, states, and current/proposed processor computation hashes.

The generic metric editor follows the same business-first rule. Metric choices are ordered and labelled as Processor · metric name · metric kind. Its Review panel leads with the metric description, a readable calculation, aggregate inputs, the metric families supported by those aggregate states, and presentation semantics. Calculations are rendered as explanatory Markdown with inline LaTeX notation; exact generated YAML remains available in collapsed technical details. Configuration Builder applies and post-validates the multi-file patch inside a rollback boundary. The changed processor requires the first ingestion run for a new workspace or replay/backfill for existing aggregates; recipe installation never starts that data operation or converts one stored sketch family into another.

Both surfaces use structured YAML parsing and the closed expression AST. Neither writes free-form Python or mutates the workspace until the user presses an explicit apply action; every apply re-runs catalog validation. Apply then classifies whether existing aggregates can open a report or whether the user must continue to Data Load. The latter is a handoff, not an implicit run.

Privacy-safe authoring instrumentation records only allowlisted workflow, stage, event, outcome, bounded duration/count, and materialization-required flags under an anonymous session journey. It has no arbitrary metadata field, so sample/field values, object identifiers, local paths, prompts, credentials, and provider error text cannot be attached. The revised entry can be hidden with VALUESTREAM_AUTHORING_V2=0 during measured rollout. See the authoring rollout guide.

11. Technology stack and rationale

Layer Technology Why
Programming language Python 3.11+ Mature data tooling; team familiarity; rich sketch/Polars/Plotly ecosystem
Vectorized execution Polars >= 1.x Lazy frames, streaming engine, expressive groupby + custom map_groups
Persistent aggregates Apache Parquet (PyArrow writer) Columnar, hive-partitioned, portable, compact
Bounded processor checkpoints Bounded rolling DuckDB state One writer/session, native exact SQL, Arrow C Stream current payload, transactional fingerprint journal
SQL surface DuckDB Read-Parquet TVF, fast aggregation, easy views, single-file metadata DBs
BI export DuckDB tables Optional materialized metric tables for Superset/SQL tools
Distribution sketches Apache DataSketches (Python datasketches) t-digest, KLL, CPC, HLL, Theta, Frequent-Items
Statistics SciPy Chi-square / G-test / odds ratios / CIs
ML helpers scikit-learn (FeatureHasher, cosine_similarity), polars_ds Personalization, weighted mean
UI framework Streamlit Quickest path to interactive dashboards in Python
Plotting Plotly Interactive charts the existing user base already knows
Read-only API FastAPI + Pydantic v2 Typed metric/chart/freshness/chat endpoints with OpenAPI
Schema validation jsonschema (Draft 2020-12) Validate YAML config against a stable schema
Config templating Jinja2 (optional, for variants) Workspace-specific overrides
Packaging uv + pyproject.toml Already used in the current repo
Quality gates Local uv commands Lint, format, type-check, test, docs build, schema-validate sample configs

Optional/future: - Apache Arrow IPC for very large in-process transfers between SDK and notebook clients. - ConnectorX / ADBC for upstream operational DB ingestion (today the upstream is files only).

12. Concurrency model

  • One pipeline run at a time per source, enforced by a file-system advisory lock at meta/source_<id>.lock.
  • A clean rebuild acquires all selected source locks in deterministic order and holds them through forced ingestion, coverage validation, scoped cleanup, and aggregate-view refresh. This prevents a concurrent run from publishing files that cleanup could mistake for superseded output.
  • Inside a run, chunks are processed sequentially by default (predictable memory profile, simpler error semantics). The --parallel <N> flag runs chunks in a process pool of N workers: partial parquet part files are per-chunk so worker writes never collide, and all ledger writes stay in the parent process (the DuckDB metadata files are single-writer). Worker processes sidestep the GIL held by Python sketch building, so the initial load scales with cores.
  • Inside a chunk, processors fan out through one batched pl.collect_all: every processor's chunk_aggregate plan collects in a single pass (sharing the scan via common-subplan elimination), with a logged sequential fallback if the batch fails.
  • source_scan remains parallel-safe because each target owns an immutable dependency frame. A source containing persistent_sharded frequency_response is different: targets run oldest-to-newest through one long-lived rolling-state writer, so requested chunk-process parallelism is capped at one for that source. Other source and processor semantics remain unchanged.
  • The query layer is read-only, fully concurrent — Streamlit, SDK, SQL export, MCP, and API reads can run while the engine is writing because Parquet writes go to immutable run-specific files. The run's durable running row is the outer publication barrier. Within it, atomic Parquet writes and complete lineage commit before the chunk's ok row, which is the chunk commit marker. A committed chunk becomes visible only after its parent run reaches any terminal state; until then readers retain the previous successful version. Publication keeps one globally latest successful attempt for each (source, chunk), so a successful empty recomputation supersedes all older output. Idempotent reuse is subject to the same rule: after a catalog rollback, an older matching contract is reprocessed and becomes the latest attempt instead of being silently republished.

13. Caching strategy

  • Reader cache (per-chunk): if the same chunk is read by multiple processors, the source LazyFrame is collected once and aliased.
  • UI metric cache (Streamlit st.cache_data, in ui/data.py): query results are memoized keyed by workspace, metric, group-by, filters, grain, date range, and a signature derived from the catalog, the processor config hash, the aggregate files' (count, mtime, size), and the ledger DBs' (mtime, size). Any ingestion run therefore invalidates the cache automatically. This cache lives in the Streamlit surface only.
  • Query layer (query/): the executor itself is intentionally stateless — SDK, MCP, and CLI callers always read live aggregates. Predicate pushdown (config-hash filter and period partition pruning) keeps the per-call cost low, and DuckDB/Parquet plus the OS page cache absorb repeated reads. A process-level LRU at this layer is a possible future addition but is deliberately not present today, so headless callers never serve stale numbers.

There is no query cache below the storage layer (no in-memory copy of the aggregate store) — Parquet + the OS page cache do that work. Processor-owned checkpoints are ingestion acceleration state, not a query-result cache; deleting them cannot change report semantics.

14. Failure semantics

Failure Detection Effect Recovery
Reader can't open a file I/O exception inside chunk loop That chunk fails; run continues with next chunk Operator fixes file; re-run the source; only failed chunks process
Processor exception Exception inside chunk_aggregate The chunk fails and none of that run's partials become query-visible; the previous successful chunk version remains visible Operator fixes config or code; re-run the source
Partial Parquet write incomplete Write done atomically (write-then-rename) Incomplete file never visible to readers None needed
Rolling processor state missing or valid-but-incompatible Metadata and transactional-journal reconciliation against the current schema-8 contract and expected raw fingerprints No aggregate is authorized by state alone; the stable database is reinitialized Rebuild bounded rolling state automatically from authoritative source files oldest-to-newest
Rolling processor state corrupt DuckDB open/recovery, schema, journal, or exact-SQL failure The affected persistent frequency source fails closed; previous successful aggregates remain visible Run with --force to rebuild state from authoritative source files, or remove rolling.duckdb and retry
Grain materialization fails Exception while deriving a configured grain The chunk remains unpublished at every grain; previous successful data remains visible Fix the cause and re-run the source
Process is terminated before run finalization A prior running row exists after the next caller acquires the source lock The interrupted run remains invisible until its committed chunks are verified The next normal source run verifies fingerprint, lineage, files, and computation hashes; valid chunks are published under a recovered partial run and reused, invalid chunks are reprocessed
Config validation fails on load JSON-Schema error Engine refuses to start; CLI prints actionable error Operator fixes YAML
Clean rebuild safety check fails Empty discovery, incomplete source run, missing published path, or catalog hash change Scoped cleanup does not start; prior aggregate files and audit metadata remain Fix discovery/config/run failure and retry the clean rebuild
Query layer error Exception during plan/execute UI/CLI surface returns a structured error; API returns a governed 4xx/5xx Operator inspects logs

A chunk is the unit of recovery; a run is the unit of reporting and the outer publication barrier. Successful chunks within a partially-failed run are kept, marked, and re-used on the next run, but the completed run itself is failed whenever any chunk fails. Data Load presents repeated chunk errors as one grouped failure with the affected chunk count and identifiers. After acquiring the source lock, a new caller treats every older running row for that source as interrupted: an ok chunk is retained only when its current input fingerprint, file lineage, physical provenance, and processor computation hashes all verify. The stale run becomes partial when at least one chunk verifies and failed otherwise. Files without a committed chunk marker remain invisible and are eligible for a later vacuum. Recovery fetches lineage once per stale run, indexes that run's physical paths once, and deep-scans schema-compatible processor/grain files together; embedded provenance verification remains mandatory. Run-level input/kept row totals sum only the chunks whose final durable marker is ok.

15. Security and privacy

Concern Posture
Source-derived identity state Not query-visible. A bounded processor may persist a minimal checkpoint; it remains sensitive, independently retained, and rebuildable from the authoritative source.
Customer hash shards Routing optimization only, not anonymization; tokenize/HMAC upstream and encrypt/control the workspace where required.
Identity sketches CPC is the distinct-count default; HLL remains supported; Theta can answer distinct count and is preferred when the same state also needs set algebra. Sketches are not a cryptographic anonymization boundary, so identifiers should be tokenized/HMACed upstream when required.
Code-injection via config No eval; only the closed expression AST
HTTP API auth Optional bearer token on loopback; CLI requires one for non-loopback binds; OIDC remains deferred
Multi-tenancy One workspace per tenant/variant; filesystem permissions enforce isolation
Audit Every aggregate row and every query response carries config_hash and lineage pointers
Logs Structured JSON; sensitive fields scrubbed before logging
Secrets Read from environment; never in YAML; .env files excluded from git
Governed SQL Disabled by default for API/MCP; when enabled, only allowlisted aggregate/export paths are accessible and DuckDB external access/extension loading is locked down

16. Observability

  • Logs (structured JSON): every chunk start/end, processor timing, row counts, memory snapshots.
  • Metrics: valuestream_chunk_seconds, valuestream_chunk_rows_in/out, valuestream_run_status, valuestream_aggregate_size_bytes, valuestream_query_seconds, valuestream_query_rows_scanned. Prometheus /metrics is reserved for the deferred service surface.
  • Tracing (OpenTelemetry, optional): one span per run, child spans per chunk, child spans per processor, attributes for config_hash, chunk_id, and group-by columns.
  • Health: the Streamlit Ops page and CLI validation expose operational health; the read-only API exposes GET /health.
  • Freshness and provenance: Streamlit exposes freshness, while API/MCP metric queries return the selected physical grain, config hashes, contributing runs/chunks, scan count, and latest aggregate timestamp.

17. Deployment topology

Value Stream runs as one process per workspace by default:

+-------------------------------------------+
|  Host (VM / container / laptop)           |
|                                           |
|   valuestream serve --workspace bdt           |
|                                           |
|   ├── Streamlit on :8501                 |
|   ├── ingestion runner (CLI / cron)      |
|   ├── optional local stdio MCP           |
|   ├── optional read-only FastAPI         |
|   └── deferred: remote HTTP MCP/OIDC     |
|                                           |
|   $WORKSPACE_DIR -> /data/valuestream/bdt    |
+-------------------------------------------+

For multiple workspaces, run multiple processes — each has its own working directory, ports, and configuration. A reverse proxy (Traefik / nginx) routes by hostname.

For local development everything runs from a single valuestream serve command pointed at a local workspace directory.

18. Lifecycle of an aggregate row

        +-------------------+    file group    +---------------------+
        |  source folder    | ---------------> | discovery & grouping|
        +-------------------+                  +---------+-----------+
                                                          |
                                                          v
                                +-------------------+--------+
                                | reader -> Polars LazyFrame |
                                +---------+------------------+
                                          |
                                          v
                              +-----------+-----------+
                              | transforms (typed AST)|
                              +-----------+-----------+
                                          |
                                          v
                              +-----------+-----------+
                              | processor.chunk_agg() |
                              +-----------+-----------+
                                          |
                                          v
                              +-----------+-----------+
                              | partial parquet write |
                              +-----------+-----------+
                                          |
                                          v
                              +-----------+-----------+
                              | (later) compaction    |
                              | daily -> monthly ->   |
                              | summary               |
                              +-----------+-----------+
                                          |
                                          v
                              +-----------+-----------+
                              | served by query layer |
                              +-----------+-----------+
                                          |
                                  (config change?)
                                          |
                                  +-------+-------+
                                  |               |
                                  v               v
                       +----------+--+   +--------+----------+
                       | re-process  |   | leave on disk     |
                       | (new run,   |   | until vacuum,     |
                       |  new hash)  |   | served as legacy  |
                       +-------------+   +-------------------+

19. Cross-cutting concerns matrix

Concern Where handled
Schema validation Config loader (jsonschema) at startup and on validate CLI
Time grains Source transforms (derive_calendar); processor base time.grain; query-time rollup planner
Sketch parameters State spec in processors.yaml (type, lg_k, k, …)
Filtering Two layers: source-wide transform filters; per-processor filter AST
Default values sources.<id>.defaults map
Renaming / casing Built-in rename_capitalize transform
Holiday/business calendars Optional calendar block under defaults; future derive_calendar extension
Reproducibility config_hash + chunk ledger
Dataset evolution Workspace migrations table tracks YAML changes over time
Retention vacuum CLI prunes superseded chunk partials and older config_hash aggregates; Data Load clean rebuild retains only the newly verified run files within its selected source scope
Backfill valuestream backfill --source X --from YYYY-MM-DD re-runs all chunks in window
Disaster recovery Aggregate store + metadata is everything; tar it up, ship it, untar

20. Boundaries — what Value Stream is not

  • Not an ETL platform — there is no transform graph, no joins of arbitrary sources, no schedule across sources.
  • Not a warehouse — there is no SELECT * on raw events.
  • Not a streaming platform — micro-batch is feasible, true streaming is out of scope.
  • Not a feature store — the score-distribution processor exists for ML monitoring, not feature serving.
  • Not a CDP — entity resolution beyond approximate distinct counting is not provided.

These boundaries keep the platform small enough to build, operate, and reason about. Anything outside the boundary is delegated to the upstream operational system or a downstream specialized tool.


Reading order for a new engineer

  1. This document.
  2. concepts/domain-model.md — concepts and their relationships.
  3. design/replacement-design.md — full DSL, schemas, APIs, migration.
  4. reference/processors.md — per-processor algorithms.
  5. reference/algorithms.md — sketches and statistical tests.
  6. reference/expression-dsl.md — formal grammar for the AST.
  7. reference/readers-and-formats.md — file format specs.
  8. reference/chart-catalog.md — Plotly chart bindings.
  9. reference/faq.md — common questions; check this when stuck.