Stream Data Archives
A Stream Data Archive is the unit of configuration and storage for time-series data in OctoMesh. Each archive is a versioned, typed, per-tenant CrateDB table that captures a curated set of attribute paths from a Construction Kit type. Archives have a strict lifecycle (Created → Activated → Disabled / Failed), a schema that is frozen once activated (the one exception is adding a declared column, see Conflict precedence), and a three-tier activation gate (instance → tenant → archive) that determines whether the data plane is open.
Archives are runtime entities — instances of the System.StreamData CK model — and can be created, activated, disabled, retried, and deleted through the Refinery Studio, GraphQL, REST, or the octo-cli.
Archive rows can be backed up too. Use
ExportArchiveData /
ImportArchiveData
to move a single archive's rows, or include all of a tenant's archives in a full tenant backup with
dump --include-archive-data — see Repository Backup & Restore.
What an Archive Is
An archive bundles three concerns into one entity:
- Target CK type — the source-of-truth for the data shape. Inheritance is supported: an archive defined on type
Aapplies to all derived types; one CrateDB table per concrete type. - Column list — the attribute paths captured at ingestion time, mapped to typed CrateDB columns. Required vs. optional and indexed vs. non-indexed are configured per column.
- Status — the lifecycle state. Inserts and queries are accepted only when status is
Activated.
Where a CK model defines the shape of data, an archive defines what slice of that data is persisted to time series. Many archives can target the same CK type — different retention, different columns, different downsampling.
Properties
| Property | Description |
|---|---|
| Per-tenant | Each tenant has its own CrateDB schema. No cross-tenant data. |
| Typed | Each captured attribute path becomes a real CrateDB column with a declared type. |
| Immutable schema | Once activated, TargetCkTypeId and the existing Columns cannot change. New version = new archive. The one exception is add-only: a newly declared ingested column is added to the table on the next activateArchive. |
| Multi-archive | A CK type may be captured by any number of archives concurrently. |
| Status-gated | Inserts and queries are accepted only while the archive's status is Activated. |
| Soft-delete preserves data | Disabling an archive blocks reads and writes but preserves the underlying CrateDB table. |
Archive Subtypes
The abstract base Archive has three concrete subtypes (System.StreamData CK model, ≥ 1.4.0):
| Subtype | Time axis | Ingestion path | Use case |
|---|---|---|---|
RawArchive | Single timestamp per row | External producers (SaveStreamDataInArchive@1, IStreamDataRepository.InsertAsync) | Instant measurements (sensor readings, machine telemetry, events) |
RollupArchive | Bucketed [bucketStart, bucketEnd) | System rollup orchestrator only — no external inserts | Derived bucketed aggregations of one or more source archives (raw, time-range, or rollup-of-rollup) |
TimeRangeArchive | Explicit [from, to) range per row | External producers (SaveTimeRangeStreamDataInArchive@1, SaveTimeRangeSeriesInArchive@1, IStreamDataRepository.InsertTimeRangeAsync) | Externally pre-aggregated data (EDA reports, smart meters, weather APIs) |
RawArchive and TimeRangeArchive are operator-defined and accept external data. RollupArchive columns are generated server-side from the configured aggregations — direct edits are rejected.
CK Model
Loaded only when StreamData is enabled at the instance level. Lives in System.StreamData (currently 1.13.1).
# Abstract base
ckType: Archive
isAbstract: true
attributes:
- TargetCkTypeId # CkId<CkTypeId>
- Columns # array<CkArchiveColumn>
- Status # CkArchiveStatus
- RawRetentionMs # int64, optional — partition-drop horizon
- MaxRetroactiveReachMs # int64, optional — bounded retro reach cap for automatic recompute (since 1.6.8); null = unbounded
- ConflictPrecedence # array<CkArchiveConflictKey>, optional — ordered keys deciding which of two writes to one row survives (since 1.13.0); empty = last write wins
ckRecord: CkArchiveConflictKey # since 1.13.0
attributes:
- Column # string — one of the archive's own columns
- Order # CkConflictKeyOrder — HigherWins | LowerWins
# Concrete subtypes
ckType: RawArchive (derivedFrom: Archive)
ckType: TimeRangeArchive (derivedFrom: Archive)
attributes:
- Period # TimeSpan, optional — advisory hint
ckType: RollupArchive (derivedFrom: Archive)
attributes:
- Sources # array<CkRollupSourceReference> — the source archives with their validity spans (since 1.8.0)
- SourceArchiveRtId # OctoObjectId, optional — DEPRECATED since 1.8.0, kept so pre-1.8.0 entities import unchanged
- BucketSizeMs # int64
- BucketAlignment # FixedSize | CalendarDay | Iso8601Week | CalendarMonth | CalendarQuarter | CalendarYear
- ReferenceTimeZone # IANA zone id (e.g. Europe/Vienna), optional — anchors calendar alignments (since 1.6.4)
- WatermarkLagMs # int64 — how far behind real time the orchestrator stays
- LastAggregatedBucketEnd # DateTime, optional
- Aggregations # array<CkRollupAggregation>
- FrozenUntil # DateTime, optional
ckEnum: CkArchiveStatus { Created, Activated, Disabled, Failed }
ckEnum: CkRollupFunction { Avg, Min, Max, Sum, Count, TimeWeightedAvg, StateDuration, First, Last }
ckEnum: BucketAlignment { FixedSize, CalendarDay, Iso8601Week, CalendarMonth, CalendarQuarter, CalendarYear }
ckRecord: CkRollupSourceReference # since 1.8.0
attributes:
- SourceArchiveRtId # String, required — the source archive (raw, time-range, or another rollup)
- ValidFrom # DateTime, optional — INCLUSIVE start; null = open start
- ValidTo # DateTime, optional — EXCLUSIVE end; null = open end
ckRecord: CkArchiveColumn
attributes:
- Path # attribute path, e.g. "sensor.reading.value" or "readings[*].value"
- Required # bool
- Indexed # bool (default true — opt-out only when justified)
ckRecord: CkRollupAggregation
attributes:
- SourcePath # attribute path on the source archive's CK type
- Function # Avg | Min | Max | Sum | Count | TimeWeightedAvg | StateDuration | First | Last
- TargetColumnName # optional override; otherwise derived from SourcePath + Function
Source declaration and the deprecated scalar
Sources replaces the single SourceArchiveRtId a rollup carried before 1.8.0. Both forms are
reconciled in exactly one place — the snapshot mapping on the read side of the runtime
store — under this rule:
A non-empty
Sourceslist wins. WhenSourcesis empty, a setSourceArchiveRtIdis read as one unbounded source. ASourceArchiveRtIdthat disagrees with a non-emptySourceslist is carried as a conflict and rejected at activation — never during enumeration, so listing a tenant's rollups can't fail because one entity is inconsistent.
The platform itself never writes SourceArchiveRtId any more: new rollups are stored with
Sources only. The scalar exists so that entities, seeds and blueprints written before 1.8.0 keep
importing unchanged, and a rollup stored in that old form behaves exactly as before — one unbounded
source is operationally identical to a pre-1.8.0 single-source rollup. No CK-model migration is
involved.
Engine decision doc: octo-construction-kit-engine/docs/concept-multi-source-rollups.md — the
authoritative, rule-level description of sources, validity spans and coverage.
Storage Layout
Naming
| Element | Convention |
|---|---|
| Schema | clean(tenantRtCkId.SemanticVersionedFullName), max 63 chars |
| Table | clean(targetCkType.SVF) + "_" + clean(archiveRtCkId.SVF), max 200 chars |
| Column | camelCase from CK attribute id; nested paths concatenated camelCase via ColumnNameMapper |
clean() strips non-alphanumerics. When a generated name exceeds the length limit, the helper truncates and appends a deterministic SHA-256 suffix ({truncated}_{hash16}). The truncation is logged at Information with both names for traceability.
Standard columns (every archive table)
| Column | Type | Notes |
|---|---|---|
rtId | TEXT | Entity instance id |
timestamp | TIMESTAMP WITH TIME ZONE | Datapoint timestamp (UTC) |
ckTypeId | TEXT | Concrete type (relevant for inheritance) |
rtCreationDateTime | TIMESTAMP WITH TIME ZONE | DEFAULT CURRENT_TIMESTAMP |
rtChangedDateTime | TIMESTAMP WITH TIME ZONE | Updated on conflict |
rtWellKnownName | TEXT | Optional |
Primary key: (timestamp, rtId, ckTypeId).
TimeRangeArchive replaces timestamp with bucketStart + bucketEnd; PK is (bucketStart, bucketEnd, rtId, ckTypeId).
Path → column mapping
| Path shape | CrateDB column type | Nullability |
|---|---|---|
Scalar (voltage) | mapped from CK attribute primitive type | NOT NULL if required |
Record (sensor) | OBJECT(STRICT) with subfields from CkRecord | as above |
Array of scalars (readings[*].value) | ARRAY(<scalar type>) | NOT NULL if required |
Array of records (readings[*]) | ARRAY(OBJECT(STRICT)) | as above |
Each generated column is indexed by CrateDB's default rules unless CkArchiveColumn.Indexed = false, in which case the DDL emits INDEX OFF. Disable indexes only for columns that are read but never filtered or aggregated on.
Upsert semantics
Inserts run as multi-row INSERT … ON CONFLICT DO UPDATE on the row key — (timestamp, rtId, ckTypeId), or (window_start, window_end, rtId, ckTypeId) on a windowed archive. By default every column of the archive is overwritten by the incoming value, so the last write wins for that row.
rtChangedDateTime is bumped on every conflict; rtCreationDateTime is never modified, and rtWellKnownName is written on insert and then left alone. Idempotent re-inserts are safe.
NULL clears the stored valueThe upsert does not merge column by column. A second producer writing the same row must carry the full set of columns it owns, or the columns it omits are nulled.
Conflict precedence — deciding which write wins
Last-write-wins means the stored value reflects the order writes arrive in, not the data. Wherever the same row can legitimately be written more than once — a corrected reading, a re-sent document, a backfill running next to a live feed — that is a correctness problem: replaying the same set of writes in a different order produces a different archive.
Archive.ConflictPrecedence (System.StreamData 1.13.0, available on every archive type) makes the update conditional. It is an ordered list of keys, each naming one of the archive's own columns plus which direction of it wins:
| Field | Meaning |
|---|---|
Column | one of the archive's own columns — the guard compares the incoming value against the stored one |
Order | HigherWins for timestamps, versions, sequence numbers · LowerWins for rank codes numbered best-first |
An empty list is the default and keeps the behaviour described above, so declaring a precedence is an explicit opt-in and changes nothing for existing archives.
The keys are compared lexicographically: the first decides, and a later one only breaks a tie in the ones before it. That makes them a total order over the data, so the surviving value is its maximum — the same value whichever write lands first. (Read as a conjunction — "better rank and newer" — the outcome would still depend on arrival order, which is the problem being solved.) This holds for writes the keys tell apart: two writes that are equal in every key are not ordered by them, and the later write replaces the stored row (an identical re-delivery stays an idempotent upsert), so declare enough keys to distinguish every pair of writes that can carry different values.
A write that loses leaves no trace at all: not the values, not rtChangedDateTime, not was_updated.
Null handling:
| Stored | Incoming | Outcome |
|---|---|---|
NULL | value | incoming wins — a row with no ranking information stays replaceable |
| value | NULL | stored wins — a write that cannot prove it is better never displaces one that can |
NULL | NULL | tie; the next key decides |
The last key accepts equality, so re-delivering an identical payload is idempotent rather than silently dropped.
Example — metering values, where an estimate must never overwrite a measurement:
ConflictPrecedence:
- { Column: DataQuality, Order: LowerWins } # 1 = measured … 3 = estimated
- { Column: SourceDocumentDate, Order: HigherWins } # the document's own timestamp
With LowerWins over an enum column, the CK enum keys decide the ranking — including values that were never meant to take part, such as an Unknown = 0 that would outrank every real reading, or a Manual correction numbered last that could then never displace anything. Check the whole enum, not just the values you had in mind.
Two constraints:
- Every key must be one of the archive's own columns; anything else is rejected when the statement is built.
- The precedence applies to ingest only. Restoring rows with
ImportArchiveDatabypasses it — an operator restoring a snapshot gets back exactly what they gave.
Adopting a precedence on an archive that is already Activated usually means declaring a new column for the ranking value. Declared-but-missing columns are added to the CrateDB table on the next activation (add-only, and a Required column is added nullable because existing rows cannot retroactively have a value) — no rows are lost, and existing rows carry NULL for the new column. Call activateArchive on the active archive to trigger it: the archive stays Activated, and if adding a column fails the error is returned while the archive keeps its status and its table. Rollup archives are left out, because their aggregate columns are derived and a new aggregation needs a backfill rather than an empty column. A precedence key has to name one of the archive's columns; an archive whose ConflictPrecedence names anything else fails to activate.
Time semantics
- Canonical timezone is UTC. Incoming timestamps are normalized; naive
DateTimevalues are interpreted as UTC. - API responses serialize as ISO-8601 with explicit
Z. - Late-arriving data is accepted via the standard upsert path.
- Future timestamps are not rejected — drift detection belongs in the pipeline / monitoring layer.
Lifecycle
| Status | Table state | Inserts / Queries | Notes |
|---|---|---|---|
Created | Not provisioned | Rejected | Fully editable. No CrateDB resources allocated. |
Activated | Provisioned, frozen | Accepted | Schema immutable except for added columns. Validation on update path. |
Disabled | Provisioned, frozen | Rejected | Data preserved. Re-enable to resume. |
Failed | May be partial | Rejected | Activation DDL failed. retryActivate is idempotent. |
Path validation runs on every transition that ends in Activated (Created/Disabled/Failed → Activated): each CkArchiveColumn.Path is re-validated against the current CK model. Failures throw ArchivePathInvalidException before any DDL runs.
Deletion is destructive: the CrateDB table is dropped (idempotent DROP TABLE IF EXISTS) and the entity is soft-deleted (rtState = Archived). Historical data is lost. Deletion is allowed from any status — except an archive of any subtype that is listed as a source of a rollup, which is rejected with RollupSourceInUseException.
Activation Layers
Every data-plane operation is gated by a three-level AND-coupled check:
Instance enabled? ── no ──▶ DI does not register CrateDB stack; tenant enable throws.
yes
▼
Tenant enabled? ── no ──▶ Archive activate throws.
yes
▼
Archive Activated? ── no ──▶ Insert / query throws ArchiveNotActivatedException.
yes
▼
OK
| Layer | Source |
|---|---|
| Instance | StreamData:Enabled in appsettings (default false) |
| Tenant | StreamDataGlobalSettings.IsEnabled (per-tenant flag) |
| Archive | CkArchive.Status == Activated |
The studio reads both flags from the aggregate tenant-features endpoint (GET {tenantId}/v1/features/status, whose streamData part carries { instanceEnabled, tenantEnabled }); if instanceEnabled = false the archive navigation is hidden. If instanceEnabled = true but tenantEnabled = false the user sees an empty state with an "Enable for tenant" action.
Authorization
Three roles, registered in the identity service alongside AdminPanelManagement:
| Role | Allowed operations |
|---|---|
StreamDataAdmin | Full lifecycle: Create / Update / Delete archives; Activate / Disable / Enable / Retry; Tenant Enable / Disable; rollup lifecycle (Freeze / Unfreeze / Rewind). |
StreamDataWriter | Insert into activated archives. Metadata read (List / Get archives). No lifecycle, no CRUD. |
StreamDataReader | Time-series queries (simple / aggregation / grouped / downsampling). Metadata read. No insert, no lifecycle. |
AdminPanelManagement implies all three. Enforcement is server-side; the studio gates buttons by route guards only. A user without write rights who clicks an action gets a notification from the server-returned ArchiveNotActivatedException / Forbidden.
Rollup Archives
A rollup archive consumes one or more source archives and writes one row per [bucketStart, bucketEnd) window. Buckets are sized by BucketSizeMs and aligned by BucketAlignment:
| Alignment | Bucket boundary |
|---|---|
FixedSize | LastAggregatedBucketEnd + N · BucketSizeMs (offset-invariant; ignores the zone) |
CalendarDay | Local midnight in ReferenceTimeZone |
Iso8601Week | Monday 00:00:00 local time in ReferenceTimeZone |
CalendarMonth | First day of month 00:00:00 local time in ReferenceTimeZone |
CalendarQuarter | 1 January / 1 April / 1 July / 1 October, 00:00:00 local time in ReferenceTimeZone (since System.StreamData 1.8.0) |
CalendarYear | January 1 00:00:00 local time in ReferenceTimeZone |
BucketSizeMs is informational only for the calendar alignments (month and quarter lengths vary). The orchestrator stays WatermarkLagMs behind real time before closing a bucket, so late-arriving raw data still lands in the right window.
CalendarQuarter follows the CalendarMonth shape exactly: alignment happens in local wall-clock time and is converted back to UTC, so DST is handled for free — in Europe/Vienna the quarter holding the spring-forward switch is 90 days minus one hour of UTC duration, the one holding the fall-back switch is 92 days plus one hour. An empty ReferenceTimeZone yields UTC quarters.
Bucket interval vs source granularity
The rule is evaluated per source, against each source's declared window length (a TimeRangeArchive's Period, or a rollup source's BucketSizeMs). Raw sources that expose no declared granularity are unconstrained. Activation rejects a violation with RollupBucketIntervalException; the check lives in RollupValidator.ValidateForActivation, so a misconfigured rollup fails fast at Activate rather than producing wrong aggregates at query time.
| Rollup alignment | Source | Accepted when |
|---|---|---|
FixedSize | any | BucketSizeMs is ≥ the source window and an integer multiple of it — e.g. over a 15-minute source, 15 min / 30 min / 1 h / 1 d are accepted, while 5 min (finer) and 20 min (not a whole multiple) are refused |
| calendar / ISO week | a calendar-aligned rollup source | the source alignment nests inside the rollup alignment: CalendarDay ⊂ CalendarMonth ⊂ CalendarQuarter ⊂ CalendarYear. Iso8601Week contains only CalendarDay and nests inside nothing but itself |
| calendar / ISO week | a time-range, fixed-size, or fixed-bucket rollup source | the source window does not exceed the alignment's longest bucket: 1 d (CalendarDay), 7 d (Iso8601Week), 31 d (CalendarMonth), 92 d (CalendarQuarter), 366 d (CalendarYear) |
The constraint prevents a single source row from being split across two output buckets, which would silently distort Sum/Count and skew Avg/Min/Max. The last row is what makes a legacy-history cutover legal: a time-range archive whose quarterly windows declare a ~92 d Period is accepted under a CalendarQuarter rollup, while anything genuinely coarser is rejected.
Validity spans
Each entry of Sources is a CkRollupSourceReference carrying the source archive plus an optional half-open validity span [ValidFrom, ValidTo) — ValidFrom inclusive, ValidTo exclusive, either bound null for an open end. A reference with neither bound is unbounded and covers every bucket. Two references that abut at the same timestamp are therefore disjoint, and the bucket that starts at the cutover belongs to the later source.
This is what lets legacy history imported at a coarse granularity and natively ingested fine-grained data share one rollup ladder:
legacy quarterly time-range archive ──[ ValidTo = 2025-10-01 ]──┐
├──► CalendarQuarter rollup ──► yearly rollup
15-minute raw archive ──► monthly rollup ──[ ValidFrom = 2025-10-01 ]──┘
Rollups stacked on top need no knowledge of the split; consumers still query one archive per request.
Validation
Every rule below names the offending source archive in its message and surfaces verbatim through the GraphQL / REST error mapping. Rows 1–6 and 12 are pure checks on the declaration and row 7 additionally walks the tenant's rollups; all of them run at create and again at activation. Rows 8–11 need each source's snapshot and run at activation only.
| # | Rule | Exception |
|---|---|---|
| 1 | At least one source is declared. | RollupSourcesRequiredException |
| 2 | Each source archive appears at most once. | DuplicateRollupSourceException |
| 3 | A bounded span is non-empty (ValidFrom < ValidTo). | RollupSourceSpanInvertedException |
| 4 | At most one source has an open start, at most one an open end. | RollupSourceSpanOpenEndConflictException |
| 5 | Spans are pairwise disjoint (abutting spans are disjoint — the intervals are half-open). | RollupSourceSpanOverlapException |
| 6 | Every span boundary lies on a bucket boundary of the rollup, in its reference time zone. A boundary aligned in UTC but not in the reference zone is rejected, and vice versa. | RollupSourceSpanNotOnBucketBoundaryException |
| 7 | No cycle in the source graph — following source edges from the rollup must never reach it again (including listing itself). | RollupSourceCycleException / RollupCycleException |
| 8 | Every source archive exists, is not soft-deleted, and is Activated. | RollupSourceMissingException / RollupSourceNotActivatedException |
| 9 | Every source targets the same CK type as the rollup. | RollupSourceTargetTypeMismatchException |
| 10 | Bucket granularity per source (table above). | RollupBucketIntervalException |
| 11 | Every aggregation path resolves on every source. | RollupSourcePathMissingException |
| 12 | A deprecated SourceArchiveRtId that disagrees with a non-empty Sources list. | RollupSourceDeclarationConflictException |
Rule 6 is what guarantees the property the whole design rests on: every bucket lies entirely within exactly one source's span, or within none.
The aggregation rules that predate multi-source support — at least one aggregation, no duplicate (SourcePath, Function) pair, a ComparisonValue on every StateDuration — now run at create as well, so createRollupArchive fails fast on inputs it used to accept silently and only reject at activation.
Per-bucket source selection
Per tick the orchestrator loads each distinct source snapshot once, then resolves the source per closed bucket:
- A span covers the bucket → the existing single-source aggregation runs unchanged against that archive, the watermark advances, the run is audited.
- No span covers the bucket (a gap bucket) → no row is written, the watermark advances past it, and the skip is logged at
Debug. Gap buckets are a declared, normal state — not an error. - The covering source is missing or not
Activated→ the tick stops with aWARNINGnaming rollup and source, and the watermark stays put. Aggregation resumes exactly where it stopped once the source is available again; nothing is silently skipped.
Because the source is fixed for a whole bucket, the result is identical to a single-source rollup over the concatenated data for Sum, Count, Min, Max, Avg, First and Last. For the LOCF-carry functions TimeWeightedAvg and StateDuration the carry-in never crosses a source boundary — the first bucket served by a later source opens without a carry from the earlier one, so its covered duration can be shorter than a concatenated single-source rollup would report. That is the only documented divergence.
Aggregation specs are logical (SourcePath + Function) and are resolved into physical columns per source: either a column the source declares verbatim (an ingested column by Path, a computed column by Name, a rollup source's physical column by its physical name), or — on a rollup source — the child aggregation with the same path and the same function, read function-preserving over the child buckets. Chained rollups written in the physical-name style (amountvalue_sum / Sum) therefore keep working unchanged.
The initial watermark of a newly activated rollup is unchanged: it is derived from the clock alone, never from a source, so a multi-source rollup starts exactly where a single-source one would.
Recompute and backfill across spans
- A retroactive write into a source propagates to the rollup clipped to that source's span on both ends. A late write into the legacy archive after the cutover, or into the native archive before it, is correctly ignored.
- A recompute of
[from, to)is split into per-source segments first and only then into chunks; each segment runs against its own source snapshot, sequentially in time order, and the job's row and window totals accumulate across segments. A range no span covers completes with zero totals. - Backfill start = the minimum over all sources of
max(earliest stored timestamp of that source, ValidFrom). A source holding no data in its span contributes nothing; if no source contributes, the backfill is a logged no-op. Backfill stays manual (rewindRollupWatermark/backfill_rollup_archive).
If imported history has the same granularity as the native data — a replay of 15-minute values into a 15-minute raw archive — there is nothing to declare here. Write it directly into the raw archive through the ordinary ingest path; the existing rollups above it pick it up, and a rewindRollupWatermark + backfill recomputes the affected buckets. Multi-source declarations are only needed when the history arrives at a coarser granularity and therefore has to live in its own archive.
Reference time zone (calendar alignments)
The calendar alignments snap each bucket to a local calendar boundary in ReferenceTimeZone — an IANA zone id such as Europe/Vienna — so the buckets are DST-correct: a CalendarDay bucket runs local-midnight to local-midnight, which is 22:00Z → 22:00Z under CEST (summer) and 23:00Z → 23:00Z under CET (winter). An empty zone falls back to UTC boundaries. FixedSize ignores the zone — a whole-hour window is offset-invariant for a whole-hour zone, so an hourly rollup needs no zone. Available since System.StreamData 1.6.4.
BucketAlignment and ReferenceTimeZone are frozen once the rollup is Activated. Switching a live rollup from FixedSize/UTC to a calendar alignment in a real zone (e.g. moving a daily rollup to CalendarDay + Europe/Vienna) requires deleting and re-creating the rollup, then re-activating and backfilling.
Aggregation functions
CkRollupFunction maps to one or two storage columns:
| Function | Stored columns | Read-time formula |
|---|---|---|
Sum | {base}_sum | as-is |
Min | {base}_min | as-is |
Max | {base}_max | as-is |
Count | {base}_count (BIGINT) | as-is |
Avg | {base}_sum, {base}_count (both stored) | sum / count — keeps chained rollups numerically correct |
First | {base}_first | value at the earliest observation in the bucket (numeric source columns) |
Last | {base}_last | value at the latest observation in the bucket (numeric source columns) |
The Columns[] projection on a rollup is generated from Aggregations[] via RollupColumnGenerator.Generate; direct edits to Columns are rejected by the validation hook.
Each declaration in Aggregations[] is independent, so one rollup archive can carry many aggregations over different attributes (e.g. Sum of energy and Max of power), plus several functions of the same attribute (e.g. Min and Max of power). All declared aggregates are produced in a single pass over each source bucket and each is exposed as its own series. To add an aggregation to a rollup that is already activated, create a new rollup with the full aggregation set and backfill it with rewindRollupWatermark — Aggregations[] is immutable after activation.
Operator controls
| Operation | Mutation | Notes |
|---|---|---|
| Create | createRollupArchive(input) | Takes sources[] (CreateRollupSourceInput: sourceArchiveRtId, optional validFrom / validTo). Server derives TargetCkTypeId and Columns from the sources + aggregations. The deprecated scalar sourceArchiveRtId input is still accepted for compatibility and is mutually exclusive with sources[]; it is mapped to one unbounded source before storing. |
| Freeze a range | freezeRollupArchive(rtId, until) | Monotonic — the new until must be ≥ the current FrozenUntil. The orchestrator stops producing buckets whose bucketEnd falls in the frozen range. |
| Unfreeze | unfreezeRollupArchive(rtId, acceptGaps) | Idempotent. acceptGaps is recorded for audit but the gap-detection guard is a follow-up. |
| Rewind the watermark | rewindRollupWatermark(rtId, toBucketEnd) | Re-aggregates from the given bucket boundary forward. Destructive: previously committed rows in the rewound range are temporarily out of sync until the orchestrator catches up. |
| Recompute a range | recomputeArchive(rtId, from, to, rtIdScope?) | Optimistic, no-downtime recompute of [from, to), optionally scoped to a single entity. Readers keep seeing a consistent snapshot throughout. See Rollup Recompute below. |
Delete guard. Any archive — raw, time-range, or rollup — that is listed as a source of a non-soft-deleted rollup cannot be deleted; deleteArchive refuses with RollupSourceInUseException, in either storage form and whatever validity span the reference carries. Remove or replace the rollup first.
createRollupArchive is the supported way to declare sources, and the studio's edit mode shows them read-only. The generic CK mutation can nevertheless change any attribute the model allows, so nothing technically prevents editing Sources on an activated rollup — and doing so is unsupported: it can put an already-populated rollup out of sync with its declaration and, if a boundary ends up off the bucket grid, make a recompute segment extend past it. Create a new rollup instead.
Disabling a source of an activated rollup stays allowed. It is not blocked and does not fail the rollup — the orchestrator simply stalls at the first bucket the disabled source serves until the source is activated again, or the watermark is rewound past it.
Rollup Recompute
When raw data changes after the fact — late values, corrections, re-ingests — the rollup rows derived from that window are stale. A recompute re-aggregates the affected windows. The hard requirement is that consumers (meshboards, pipelines, the API) keep working while a recompute runs: no 500, no aborted pipeline, and every read returns a consistent snapshot — either the previous values or the new ones, never a half-written mix.
For the engineer-facing internals — how a retroactive write is detected (RetroactiveWriteDetector, the consumed-watermark criterion), the dirty-windows ledger, the dependency-graph propagation, the periodic orchestrator, the optimistic atomic swap with an end-to-end sequence diagram, and a comparison of all four (re)compute paths — see Rollup Recompute Internals.
Triggers
| Trigger | Surface | Behaviour |
|---|---|---|
| Manual | recomputeArchive GraphQL mutation · POST /archives/{rollupRtId}/recompute · octo-cli RecomputeArchive | An operator forces a recompute over [from, to), optionally scoped to one entity (rtIdScope). Requires StreamDataAdmin. |
| Periodic | RecomputeOrchestratorHostedService (per-tenant tick) | Drains the dirty-windows ledger and recomputes only the windows actually marked dirty — not the whole archive. |
| Chain propagation | internal | After a rollup's recompute commits, its own dependents are marked dirty so a multi-level rollup-of-rollup chain rolls forward top-down. |
The periodic trigger is the automatic path — it drains the dirty-windows ledger the detector fills on every retroactive write, so no operator has to reset a threshold by hand. A single very-late write can only drag an automatic recompute back as far as the per-archive bounded retro reach cap (Archive.MaxRetroactiveReachMs, since 1.6.8; null = unbounded), combined with the host StreamData:Recompute:MaxRetroactiveReachHardLimitMs ceiling. Manual recomputeArchive and rewindRollupWatermark stay unbounded — the deliberate escape hatch for a deep correction (AB#4196).
Optimistic atomic swap (how consistency is guaranteed)
CrateDB has no multi-statement transaction, so the swap uses a per-window generation pointer:
- Each rollup table carries a
generationcolumn (BIGINT NOT NULL DEFAULT 0) that is part of the primary key. Normal forward aggregation always writesgeneration = 0. - A recompute aggregates the range into a per-archive staging table, then copies those rows into the live table stamped with the next generation
N+1. The previous generation stays in place and visible. - A tiny per-rollup side-table (
archive_<rtId>__genmap) holds the active generation per range. Flipping that pointer toN+1is a single-row write — the atomic commit point. - The read path injects
generation = CASE WHEN <range> THEN <activeGen> … ELSE 0 END, so every query resolves exactly one consistent generation per window. - After the flip a background sweep deletes the now-superseded generation's rows.
Failure semantics: a job that crashes before the flip leaves the staging table orphaned (reclaimed by name on the next run) and the live table + pointer untouched — readers keep seeing the previous values. A crash after the flip but before the sweep only leaves dead rows the next sweep reclaims. No partial state is ever observable.
Scoping
rtIdScope restricts the whole pipeline — staging aggregation, the generation pointer, and the post-flip sweep — to a single entity (metering point / stream). A late value for one stream therefore recomputes only that stream's windows and never touches unrelated entities sharing the same bucket. Omitting rtIdScope recomputes the entire range.
Coalescing
Overlapping triggers for the same archive coalesce: their ranges merge into one active job, and the superseded trigger is recorded with state Coalesced in the job history. There is at most one active recompute per archive at a time.
Interaction with rewind
rewindRollupWatermark re-aggregates a range forward at generation = 0. If that range had previously been recomputed, the active-generation pointer would otherwise keep readers pinned to the stale recomputed generation. The rewind therefore reconciles the generation map: it drops the pointers reaching past the rewind boundary and deletes the orphaned generation > 0 rows there, so the forward re-aggregation becomes authoritative again. This is a no-op when stream data is disabled or nothing was recomputed.
Job history & observability
Every recompute is recorded as a job snapshot, queryable newest-first (capped at 50):
- GraphQL —
runtime.streamData.recomputeJobsFor(rtId)returns[RecomputeJobInfo!]!. - REST —
GET /archives/{archiveRtId}/recompute-jobs. - octo-cli —
ListRecomputeJobs.
RecomputeJobInfo carries state (Running / Completed / Failed / Coalesced), rowsProcessed, windowsProcessed, startedAt, finishedAt, durationMs, and errorReason (set on failure). In addition, the per-rollup RollupArchiveInfo (via rollupsFor) exposes live health fields: recomputeInProgress, lastRecomputeStartedAt, lastRecomputeSuccessAt, lastRecomputeFailureAt, lastRecomputeFailureReason, dirtyWindowsPending, and pendingRecomputeRanges. Failures are also emitted through the archive audit trail.
Archive Coverage
Coverage answers "which resolution is actually available for which time range?" — the question a chart has to answer before it can decide which rung of a rollup ladder to read. Since System.StreamData 1.8.0 every archive, base and rollup, exposes it.
Measurement and timestamp convention
Coverage is measured from CrateDB on request, never derived from the model or the validity spans:
| Archive shape | availableFrom | availableTo |
|---|---|---|
| windowed (rollup, time-range) | MIN(window_start) — the start of the earliest stored window | MAX(window_end) — the end of the latest stored window (the exclusive end of the last bucket, not its start) |
| raw | MIN(timestamp) | MAX(timestamp) |
- Coverage is archive-wide, not per entity
rtId, and is a single from/to pair. - Null semantics — an archive without rows, or without a backing table, has no coverage: both fields are
null. Never a sentinel range, never an error. Every other storage error propagates. - Coverage does not show gaps. Holes inside the range are deliberately not represented, so a family whose native ingestion started long after the legacy history ended still reports one contiguous-looking range.
Caching
Answers are memoised per (tenant, archive) with a default TTL of 60 seconds, configurable per host via StreamData:Coverage:CacheTtlSeconds (mesh-adapter environment variable StreamDataCoverageCacheTtlSeconds). A TTL of zero disables memoisation; a negative TTL is rejected. Entries are invalidated on an archive delete, on completion of a recompute or backfill job, and tenant-wide when a tenant's stream data is dropped or the feature is disabled. Ordinary ingest is deliberately not invalidated — it only extends the range and the TTL absorbs it. Within the TTL a query issued seconds after data arrived can therefore still route to a coarser rung; lower the TTL if a deployment needs it fresher.
The family query
One call returns the queried archive plus every rollup transitively reachable from it, in breadth-first order with the queried archive first. It accepts any archive rtId of the family:
| Surface | Call |
|---|---|
| GraphQL | streamData.coverageFor(rtId) |
| REST | GET /api/v1/streamData/archives/{archiveRtId}/coverage |
| MCP | get_archive_coverage |
Each rung carries archiveRtId, rtWellKnownName, status, isBase (the rung is not a rollup), bucketSizeMs (rollup: BucketSize; time-range base: Period; raw base: null), bucketAlignment, storedFunctions (empty for a base archive), and availableFrom / availableTo. An unknown rtId, or a tenant without stream data, yields an empty list, not an error.
A multi-source rollup belongs to the family of each of its sources and reports the same measured coverage in each — listed exactly once per family.
Coverage is measured per archive, not derived from the graph. A disabled rung keeps reporting the range of the rows it still holds while no new data arrives, so its availableTo can look current although the rung stopped advancing. Read status from the family query alongside the coverage to tell a disabled rung from a live one.
Coverage in resolveSeriesQuery
The resolution-aware resolver applies coverage as a filter over the candidate rungs before its existing selection rule:
- A rung covers the request when its
availableFromis at or before the requested start (equal counts as covering; the requested end is not considered). A rung without coverage never covers. - The base rung takes part in the filter like every other rung.
- If no rung reports coverage at all, the filter is inert — the plan is byte-for-byte the pre-coverage result, including its signal. Hosts that do not wire a coverage provider keep the old behaviour.
- Candidates are the covering rungs; if none covers, the candidates are the rungs with the earliest
availableFrom, ties kept. - The existing selection rules then run unchanged over the candidates.
When the filter changes the chosen archive the result carries signal = CoverageLimited — which takes precedence over Ok and ResolutionLimited — plus actualPoints (the delivered point count), a diagnostic naming the excluded finer rung and its available-from, and the first-class field finerRungAvailableFrom (null when the excluded rung reports no coverage at all).
Two deliberate edge cases: if the filtered plan lands on a refusal (NoSuitableRollup / UnknownBaseGrain) that signal is kept and the coverage exclusion is prepended to the diagnostic; and if the covering candidates would yield EmptyLadder, the unfiltered result is returned — so a CoverageLimited answer always names a covering fallback rung.
GetQueryById@1 routes through the resolver and, as before, keeps the persisted archive with a warning for any signal other than Ok, routing only through the existing bucket-grid gate.
Time-Range Archives
For externally pre-aggregated data with explicit [from, to) windows. There is no source archive and no orchestrator on this path; producers push rows directly:
- GraphQL mutation:
createTimeRangeArchive(input)— same admin-grade role guard as rollup create. - Pipeline nodes:
SaveTimeRangeStreamDataInArchive@1for rows that already carry their sourcertId,SaveTimeRangeSeriesInArchive@1for whole series (one anchor entity per series). - REST:
POST /api/v1/streamData/archives/{archiveRtId}/insertTimeRange.
Period on TimeRangeArchive is advisory only — it lets the studio render a default time range but is not enforced at insert time.
Computed Columns
A computed column on a raw or time-range archive is derived by a formula from one or more other columns of the same row, rather than supplied by the producer — e.g. powerFactor = activePower / apparentPower. The value is materialised as a real, typed CrateDB column, so it is natively aggregatable by rollups and visible to direct SQL consumers (Grafana).
- Where it is evaluated — server-side in the CrateDB repository on every ingest path (pipeline node, REST, SDK), so the value is identical regardless of producer. The formula uses the shared mXparser dialect (see the Formula Expressions reference under Communication).
- Result type — declared explicitly and cast back from the
doublemXparser computes:Boolean,Int,Int64,Double, orDateTime(ticks).Stringand non-scalar types are not supported. A formula / NaN / NULL-operand error storesNULLfor that cell; ingest never fails. - CK model —
CkArchiveColumncarriesName+Formula+ResultType+ the engine-managedComputedState(sinceSystem.StreamData1.5.0) andComputedVersion/PendingFormula(since 1.6.1 / 1.6.2 respectively). A column with aFormulais computed; one with aPathis ingested. The two are mutually exclusive.
Active-archive lifecycle (optimistic / atomic)
Adding, re-formulating, or removing a computed column works on an Activated archive without blocking ingest or breaking readers, reusing the optimistic semantics of Rollup Recompute:
- Add — the column is persisted
Pending, the physical column is added (ALTER TABLE ADD COLUMN), a backfill populates it across the existing rows while it stays hidden from readers, then it flips toActiveatomically. A backfill failure leaves itFailedand the previous archive state intact. - Change formula — the column keeps serving its previous values (version N) while the new formula is backfilled into a versioned physical column
{base}__v{N+1}; ingest dual-writes both during the backfill; then a single store write swaps the version pointer atomically. A failure reverts to the previous formula. - Remove — dropped from the logical projection; the physical column is left as a harmless nullable orphan the read path no longer projects.
Limitation (formula change). Because a formula change re-versions the column's physical name, changing the formula of a computed column that another computed column references is rejected (ComputedColumnInvalidException) — re-point or remove the dependent first. Direct-SQL / Grafana consumers that bind to the physical name must likewise follow the rename.
Rollups over computed columns
A CkRollupAggregation.SourcePath may target a source computed column by its Name — the rollup aggregates it natively like any ingested column (single- and multi-aggregate). A RollupArchive may also carry its own computed columns over its aggregated output columns.
API Surface
GraphQL — lifecycle mutations
Under mutation.streamData:
| Mutation | Returns | Description |
|---|---|---|
activateArchive(rtId) | ArchiveTransitionResult | Provisions the CrateDB table, transitions to Activated. |
disableArchive(rtId) | ArchiveTransitionResult | Activated → Disabled. Table is preserved. |
enableArchive(rtId) | ArchiveTransitionResult | Disabled → Activated. Re-validates column paths. |
retryArchiveActivation(rtId) | ArchiveTransitionResult | Failed → Activated. Idempotent (CREATE TABLE IF NOT EXISTS). |
deleteArchive(rtId) | Boolean | Drops the CrateDB table and soft-deletes the entity. Rejected if any rollup lists it as a source. |
createRollupArchive(input) | OctoObjectId | sources[] + aggregations; server-derived TargetCkTypeId + Columns. |
createTimeRangeArchive(input) | OctoObjectId | Operator-defined target type and columns. |
freezeRollupArchive(rtId, until) | ArchiveTransitionResult | Monotonic FrozenUntil set. |
unfreezeRollupArchive(rtId, acceptGaps) | ArchiveTransitionResult | Clears FrozenUntil. |
rewindRollupWatermark(rtId, toBucketEnd) | ArchiveTransitionResult | Re-aggregates from the given boundary. |
recomputeArchive(rtId, from, to, rtIdScope?) | RecomputeJobInfo | Triggers/coalesces an optimistic recompute of [from, to), optionally per entity. Returns the job snapshot. |
addComputedColumn(rtId, name, formula, resultType, indexed?) | ArchiveTransitionResult | Adds a computed column and backfills it (optimistic / atomic). |
updateComputedColumnFormula(rtId, name, formula) | ArchiveTransitionResult | Re-formulates a computed column (versioned backfill + atomic swap). Rejected when another computed column references it. |
removeComputedColumn(rtId, name) | ArchiveTransitionResult | Removes a computed column (physical column left as an orphan). |
The matching read queries are runtime.streamData.recomputeJobsFor(rtId): [RecomputeJobInfo!]! (most recent jobs, newest first, capped at 50) and streamData.coverageFor(rtId) (see Archive Coverage).
Rollup read DTOs (rollupsFor, the generated CRUD projections) expose sources[] with the validity spans, plus the deprecated derived field sourceArchiveRtId — populated only when the rollup declares exactly one unbounded source, and null for every genuinely multi-source rollup and for a single source that carries a span. Clients that render the source themselves apply the same precedence for display: a non-empty sources wins, otherwise the deprecated scalar is shown as one unbounded source.
All mutations require StreamDataAdmin. Domain exceptions (§Exceptions below) surface as stable GraphQL error codes via error.extensions.OctoDetails.
CRUD on archives (List / Get / Update) is auto-generated from the CK model and lives under runtime.systemStreamDataArchives / …RawArchives / …RollupArchives / …TimeRangeArchives. The validation hook enforces immutability after activation.
GraphQL — path enumeration
query availableArchivePaths(ckTypeId: ID!, maxDepth: Int = 5): [ArchivePathInfo!]!
type ArchivePathInfo {
path: String! # "sensor.reading.value", "readings[*].value"
primitiveType: String # null when isRecord = true
isRecord: Boolean!
isArray: Boolean!
recordTypeId: String
inheritedFromCkTypeId: String # null if attribute belongs to ckTypeId itself
}
The studio's AttributePathPicker consumes this for archive create/edit. Admin-only.
REST — /api/v1/streamData
| Method | Path | Description |
|---|---|---|
| POST | /enable / /disable | Tenant-level toggle. /disable answers 409 while any archive is still Activated (AB#4255). |
| POST | /archives/{archiveRtId}/activate | Same as GraphQL activateArchive. |
| POST | /archives/{archiveRtId}/disable / /enable / /retry | Status transitions. |
| DELETE | /archives/{archiveRtId} | Same as GraphQL deleteArchive. |
| POST | /archives/{rollupRtId}/freeze / /unfreeze / /rewind | Rollup-only controls. |
| POST | /archives/{rollupRtId}/recompute | Optimistic recompute of [from, to) (rtIdScope optional). Returns the job snapshot. |
| GET | /archives/{archiveRtId}/recompute-jobs | Recent recompute jobs (newest first) for debugging. |
| POST | /archives/{archiveRtId}/insertTimeRange | Bulk time-range insert. |
| GET | /archives/{archiveRtId}/rollups | Active rollups referencing this archive as a source. |
| GET | /archives/{archiveRtId}/coverage | Measured coverage of the archive's whole family (see Archive Coverage). |
The enabled state is read through the aggregate GET {tenantId}/v1/features/status (AB#4884, replacing the
former GET /api/v1/streamData/status): it reports all four tenant capabilities the tenant delete/detach guard evaluates —
streamData with { instanceEnabled, tenantEnabled }, plus communication, reporting and aiServices
with { tenantEnabled } each. A read failure answers 500; the state is never reported as "disabled" when it
cannot be read.
Pipeline node — SaveStreamDataInArchive@1
A PathNodeConfiguration node with a required property ArchiveRtId : OctoObjectId. The node routes each incoming StreamDataPoint to the configured archive via IStreamDataRepository.InsertAsync.
Two companions write to a TimeRangeArchive:
SaveTimeRangeStreamDataInArchive@1— takes rows that already carry the sourcertId, reads the window from two attribute paths, and verifies that every distinct sourcertIdin the batch exists before inserting.SaveTimeRangeSeriesInArchive@1— takes a whole series (a key plus its array of windowed values) and does both halves itself: it resolves or creates one anchor entity per series, then bulk-inserts every value as an archive row. Use it when the alternative would be building a runtime entity per window merely to satisfy the node above — a real ingest measured 78 rows/s that way, against 1,264 rows/s here. Before it writes anything the node checks that the archive exists, is a time-range archive, isActivatedand targets the configured CK type, so a misconfigured run fails instead of leaving anchors without measurements. A value whose window is missing or does not end after it starts is skipped, and a series without a single usable window gets no anchor.
The mesh-adapter service account holds StreamDataWriter. Targeting an inactive archive surfaces a clear pipeline error (ArchiveNotActivatedException).
SDK — IStreamDataRepository
// Data plane
Task InsertAsync(OctoObjectId archiveRtId, StreamDataPoint point);
Task InsertAsync(OctoObjectId archiveRtId, IEnumerable<StreamDataPoint> points);
Task InsertTimeRangeAsync(OctoObjectId archiveRtId, IEnumerable<TimeRangeStreamDataPoint> points);
Task<StreamDataQueryResult> ExecuteQueryAsync(OctoObjectId archiveRtId, StreamDataQueryOptions options);
Task<StreamDataQueryResult> ExecuteAggregationQueryAsync(OctoObjectId archiveRtId, StreamDataAggregationQueryOptions options);
Task<StreamDataQueryResult> ExecuteGroupedAggregationQueryAsync(OctoObjectId archiveRtId, StreamDataGroupedAggregationQueryOptions options);
Task<StreamDataQueryResult> ExecuteDownsamplingQueryAsync(OctoObjectId archiveRtId, StreamDataDownsamplingQueryOptions options);
// Control plane (called by IArchiveLifecycleService)
Task EnsureArchiveCreatedAsync(CkArchive archive);
Task DeleteArchiveAsync(CkArchive archive);
Bulk inserts use all-or-nothing pre-validation: every point is validated before any SQL is sent; a violation throws RequiredAttributeMissingException carrying the index of the offending point — no row is written. A DB-level failure after pre-validation surfaces as exception.
Bulk Insert Semantics
- Resolve archive (single lookup), status check.
- Pre-validate every point in the batch — required-path coverage, type compatibility, path resolution. Stops at the first violation.
- On violation → throw with the offending index. No row written.
- On success → single multi-row
INSERT … ON CONFLICT …execution.
Pipeline callers see batch success or full failure; partial-commit scenarios are limited to transient DB failures (operational, see Reconciliation below).
Retention
Retention is partition-based and per-archive (work in progress — design captured here; some attributes are not yet exposed in the create form).
- DDL: archive tables are
PARTITIONED BY DATE_TRUNC('day', timestamp). Drop-partition is metadata-only on CrateDB. RawRetentionMson theArchivebase type defines the retention horizon.null= retain forever.- The
ArchiveRetentionScheduler(IHostedService, hourly) drops partitions older thannow() - RawRetentionMsper archive, per tenant. Failures are isolated; one bad archive does not block the rest.
For downsampling-style retention (raw rows folded into a coarser archive before deletion), use a RollupArchive pointing at the source archive — the rollup is the long-term store while the raw archive's retention can be aggressive.
Reconciliation
CkArchive entities live in MongoDB; archive tables live in CrateDB. There is no shared transaction. The lifecycle service keeps the two reconcilable without 2PC:
- All DDL is idempotent (
IF EXISTS/IF NOT EXISTS). - Operation order is Crate-first, Mongo-last:
ActivaterunsCREATE TABLE, then flipsStatusin Mongo.DeleterunsDROP TABLE, then soft-deletes the entity. Mongo is the source of truth. - Every data-plane call checks
Status == Activatedbefore touching Crate. While step 1 has run but step 2 has not, the API rejects withArchiveNotActivatedException. - On startup the
ArchiveReconcilerenumeratesActivatedarchives per tenant and (re-)provisions any missing CrateDB tables. Drift in the other direction (table without entity) is logged but not auto-dropped.
Concurrency
| Scenario | Behavior |
|---|---|
Two parallel Activate on the same archive | Serialized via RepositoryDistributedLockService. The second call observes Activated and returns idempotently. |
Disable while inserts are in flight | In-flight inserts complete. New inserts after the flip raise ArchiveNotActivatedException. |
Delete while inserts are in flight | Same as Disable, plus the DROP TABLE may retry once with a short backoff if the table is briefly locked. |
| Activate fails mid-DDL | Status → Failed. EnsureArchiveCreatedAsync is idempotent (CREATE TABLE IF NOT EXISTS), so retry is safe. |
| Bulk insert on disabled archive | Whole batch rejected on the single status lookup. |
Exception Catalog
All in Meshmakers.Octo.Runtime.Contracts.StreamData. Common base: StreamDataException (extends RuntimeRepositoryException). Surfaced via GraphQL error.extensions.OctoDetails and via REST status codes.
| Exception | Thrown when |
|---|---|
StreamDataException | Base — generic StreamData failure. |
StreamDataNotEnabledException | Instance or tenant flag is false on an operation that requires it. |
ArchiveNotFoundException | archiveRtId does not resolve to an Archive entity. |
ArchiveNotActivatedException | Archive exists but Status ≠ Activated on insert / query. |
ArchiveSchemaImmutableException | Post-activation update tries to change schema-relevant fields. |
ArchivePathInvalidException | A CkArchiveColumn.Path cannot be resolved against TargetCkTypeId. |
ArchiveColumnTypeUnsupportedException | An attribute type cannot be mapped to a CrateDB column type. |
RequiredAttributeMissingException | Insert lacks a value for a required: true scalar path. Carries the offending point index. |
ArchiveActivationFailedException | DDL execution failed during activation. Wraps the underlying SQL error. |
InvalidArchiveStateTransitionException | Lifecycle mutation called against an incompatible status. |
RollupSourceInUseException | deleteArchive invoked on an archive that is still listed as a source of a rollup. |
RollupSourceMissingException | A declared source archiveRtId does not exist or is soft-deleted. |
RollupSourceNotActivatedException | A declared source archive exists but is not Activated. |
RollupSourcesRequiredException | A rollup declares no source at all. |
DuplicateRollupSourceException | The same source archive is listed more than once. |
RollupSourceSpanInvertedException | A bounded validity span has ValidFrom >= ValidTo. |
RollupSourceSpanOpenEndConflictException | More than one source has an open start, or more than one an open end. |
RollupSourceSpanOverlapException | Two validity spans overlap. |
RollupSourceSpanNotOnBucketBoundaryException | A span boundary is not on the rollup's bucket grid in its reference time zone. |
RollupSourceTargetTypeMismatchException | A source targets a different CK type than the rollup. |
RollupSourcePathMissingException | An aggregation path/function pair does not resolve on one of the sources. |
RollupSourceCycleException / RollupCycleException | The declared sources form a cycle, or the rollup lists itself. |
RollupSourceDeclarationConflictException | The deprecated SourceArchiveRtId disagrees with a non-empty Sources list. |
RollupBucketIntervalException | A source's granularity is incompatible with the rollup's bucket size / alignment. |
Observability
Metrics
| Metric | Type | Tags |
|---|---|---|
streamdata.archive.count | Gauge | tenant, status |
streamdata.archive.status_transitions | Counter | tenant, archive, from, to |
streamdata.archive.activation_duration_ms | Histogram | tenant, archive, outcome |
streamdata.insert.duration_ms | Histogram | tenant, archive, batch_size_bucket |
streamdata.insert.points | Counter | tenant, archive |
streamdata.insert.required_violations | Counter | tenant, archive |
streamdata.query.duration_ms | Histogram | tenant, archive, query_type |
streamdata.query.rows_returned | Histogram | tenant, archive, query_type |
streamdata.crate.connections.open | Gauge | (instance-wide) |
streamdata.crate.retention_partitions_dropped | Counter | tenant, archive |
streamdata.crate.rollup_rows_inserted | Counter | tenant, rollup_archive |
Traces
Activity sources: Meshmakers.Octo.StreamData (top-level operations) and Meshmakers.Octo.StreamData.Crate (DB calls). Span attributes: streamdata.archive.rtid, streamdata.archive.target_cktype, streamdata.tenant.
Audit trail
Status transitions are emitted as RtEvent records through IEventRepository with source = RtEventSourcesEnum.StreamData. The triggering user is encoded in the event message; auto-transitions use system:reconciliation. The studio's archive details view exposes the history via the standard event query path.
Creating an Archive
From the CLI
# Tenant-level enable (required once per tenant)
octo-cli -c EnableStreamData
# Archives are runtime entities — create them through the studio, GraphQL,
# or by rt-importing a runtime-model YAML. Then:
octo-cli -c ActivateArchive -id 69fda707d47638c68edc7fea
See octo-cli — ActivateArchive and related archive commands in the sidebar for the full command surface.
From the Studio
Navigate to Repository → Stream Data Archives. The list shows every archive on the active tenant with a status badge and per-row actions (Activate, Disable, Enable, Retry, Delete, Copy ID). The New Archive, New Rollup, and New Time-Range entry points walk through the create flow with the AttributePathPicker for column selection.
See Studio: Archives for the UI walkthrough.
From GraphQL
# 1. Create the archive via the auto-generated runtime mutation
mutation {
runtime {
systemStreamDataRawArchives {
create(input: {
rtWellKnownName: "EnergyMeterArchive",
targetCkTypeId: "Industry.Energy/EnergyMeter",
columns: [
{ path: "voltage", required: true, indexed: true },
{ path: "current", required: true, indexed: true },
{ path: "powerFactor", required: false, indexed: false }
]
}) { item { rtId } }
}
}
}
# 2. Activate
mutation {
streamData {
activateArchive(rtId: "65d5c447b420da3fb1238201") {
rtId
status
}
}
}
Best Practices
- One archive per concern. Different retention, different query patterns, different consumers → different archives.
- Be explicit about required. Required scalars become
NOT NULLcolumns and are validated app-side before insert. - Index sparingly. Indexes cost storage and insert throughput. Default to indexed; opt out for columns that are read but never filtered.
- Roll up early. A
RollupArchiveis far cheaper to query over long horizons than a raw scan. - Pick alignment for the use case. Energy / EDA workflows want calendar alignments (monthly, weekly); engineering telemetry usually wants
FixedSize. - Freeze before bulk imports. When backfilling a source archive's history, freeze any dependent rollups, finish the import, then unfreeze with
acceptGaps: false(or rewind the watermark) to re-aggregate cleanly. - Soft-delete first. Disable an archive to confirm no consumer breaks before deleting.
See Also
- Rollups & recompute — detection criterion, dirty ledger, dependency graph, optimistic atomic swap
- Stream Data Access (Overview) — query surfaces (typed / transient / persisted)
- Simple Query — raw time-series retrieval
- Aggregation Queries — server-side aggregation
- Downsampling Query — bucketed downsampling
- octo-cli — EnableStreamData and related archive commands (Activate / Disable / Enable / Retry / Delete) — CLI commands for archive operations
- Studio: Archives — UI walkthrough in the Refinery Studio