Skip to main content

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.

Backing up archive data

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 A applies 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​

PropertyDescription
Per-tenantEach tenant has its own CrateDB schema. No cross-tenant data.
TypedEach captured attribute path becomes a real CrateDB column with a declared type.
Immutable schemaOnce 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-archiveA CK type may be captured by any number of archives concurrently.
Status-gatedInserts and queries are accepted only while the archive's status is Activated.
Soft-delete preserves dataDisabling 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):

SubtypeTime axisIngestion pathUse case
RawArchiveSingle timestamp per rowExternal producers (SaveStreamDataInArchive@1, IStreamDataRepository.InsertAsync)Instant measurements (sensor readings, machine telemetry, events)
RollupArchiveBucketed [bucketStart, bucketEnd)System rollup orchestrator only — no external insertsDerived bucketed aggregations of one or more source archives (raw, time-range, or rollup-of-rollup)
TimeRangeArchiveExplicit [from, to) range per rowExternal 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 Sources list wins. When Sources is empty, a set SourceArchiveRtId is read as one unbounded source. A SourceArchiveRtId that disagrees with a non-empty Sources list 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​

ElementConvention
Schemaclean(tenantRtCkId.SemanticVersionedFullName), max 63 chars
Tableclean(targetCkType.SVF) + "_" + clean(archiveRtCkId.SVF), max 200 chars
ColumncamelCase 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)​

ColumnTypeNotes
rtIdTEXTEntity instance id
timestampTIMESTAMP WITH TIME ZONEDatapoint timestamp (UTC)
ckTypeIdTEXTConcrete type (relevant for inheritance)
rtCreationDateTimeTIMESTAMP WITH TIME ZONEDEFAULT CURRENT_TIMESTAMP
rtChangedDateTimeTIMESTAMP WITH TIME ZONEUpdated on conflict
rtWellKnownNameTEXTOptional

Primary key: (timestamp, rtId, ckTypeId).

TimeRangeArchive replaces timestamp with bucketStart + bucketEnd; PK is (bucketStart, bucketEnd, rtId, ckTypeId).

Path → column mapping​

Path shapeCrateDB column typeNullability
Scalar (voltage)mapped from CK attribute primitive typeNOT NULL if required
Record (sensor)OBJECT(STRICT) with subfields from CkRecordas 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.

An incoming NULL clears the stored value

The 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:

FieldMeaning
Columnone of the archive's own columns — the guard compares the incoming value against the stored one
OrderHigherWins 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:

StoredIncomingOutcome
NULLvalueincoming wins — a row with no ranking information stays replaceable
valueNULLstored wins — a write that cannot prove it is better never displaces one that can
NULLNULLtie; 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
The enum's numbering is the precedence

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 ImportArchiveData bypasses 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 DateTime values 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​

StatusTable stateInserts / QueriesNotes
CreatedNot provisionedRejectedFully editable. No CrateDB resources allocated.
ActivatedProvisioned, frozenAcceptedSchema immutable except for added columns. Validation on update path.
DisabledProvisioned, frozenRejectedData preserved. Re-enable to resume.
FailedMay be partialRejectedActivation 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
LayerSource
InstanceStreamData:Enabled in appsettings (default false)
TenantStreamDataGlobalSettings.IsEnabled (per-tenant flag)
ArchiveCkArchive.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:

RoleAllowed operations
StreamDataAdminFull lifecycle: Create / Update / Delete archives; Activate / Disable / Enable / Retry; Tenant Enable / Disable; rollup lifecycle (Freeze / Unfreeze / Rewind).
StreamDataWriterInsert into activated archives. Metadata read (List / Get archives). No lifecycle, no CRUD.
StreamDataReaderTime-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:

AlignmentBucket boundary
FixedSizeLastAggregatedBucketEnd + N · BucketSizeMs (offset-invariant; ignores the zone)
CalendarDayLocal midnight in ReferenceTimeZone
Iso8601WeekMonday 00:00:00 local time in ReferenceTimeZone
CalendarMonthFirst day of month 00:00:00 local time in ReferenceTimeZone
CalendarQuarter1 January / 1 April / 1 July / 1 October, 00:00:00 local time in ReferenceTimeZone (since System.StreamData 1.8.0)
CalendarYearJanuary 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 alignmentSourceAccepted when
FixedSizeanyBucketSizeMs 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 weeka calendar-aligned rollup sourcethe source alignment nests inside the rollup alignment: CalendarDay ⊂ CalendarMonth ⊂ CalendarQuarter ⊂ CalendarYear. Iso8601Week contains only CalendarDay and nests inside nothing but itself
calendar / ISO weeka time-range, fixed-size, or fixed-bucket rollup sourcethe 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.

#RuleException
1At least one source is declared.RollupSourcesRequiredException
2Each source archive appears at most once.DuplicateRollupSourceException
3A bounded span is non-empty (ValidFrom < ValidTo).RollupSourceSpanInvertedException
4At most one source has an open start, at most one an open end.RollupSourceSpanOpenEndConflictException
5Spans are pairwise disjoint (abutting spans are disjoint — the intervals are half-open).RollupSourceSpanOverlapException
6Every 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
7No cycle in the source graph — following source edges from the rollup must never reach it again (including listing itself).RollupSourceCycleException / RollupCycleException
8Every source archive exists, is not soft-deleted, and is Activated.RollupSourceMissingException / RollupSourceNotActivatedException
9Every source targets the same CK type as the rollup.RollupSourceTargetTypeMismatchException
10Bucket granularity per source (table above).RollupBucketIntervalException
11Every aggregation path resolves on every source.RollupSourcePathMissingException
12A 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 a WARNING naming 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).
History at the same granularity needs no rollup change

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.

Alignment and zone are immutable after activation

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:

FunctionStored columnsRead-time formula
Sum{base}_sumas-is
Min{base}_minas-is
Max{base}_maxas-is
Count{base}_count (BIGINT)as-is
Avg{base}_sum, {base}_count (both stored)sum / count — keeps chained rollups numerically correct
First{base}_firstvalue at the earliest observation in the bucket (numeric source columns)
Last{base}_lastvalue 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​

OperationMutationNotes
CreatecreateRollupArchive(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 rangefreezeRollupArchive(rtId, until)Monotonic — the new until must be ≥ the current FrozenUntil. The orchestrator stops producing buckets whose bucketEnd falls in the frozen range.
UnfreezeunfreezeRollupArchive(rtId, acceptGaps)Idempotent. acceptGaps is recorded for audit but the gap-detection guard is a follow-up.
Rewind the watermarkrewindRollupWatermark(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 rangerecomputeArchive(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.

Sources are immutable after activation — documented, not enforced

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.

Deep dive

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​

TriggerSurfaceBehaviour
ManualrecomputeArchive GraphQL mutation · POST /archives/{rollupRtId}/recompute · octo-cli RecomputeArchiveAn operator forces a recompute over [from, to), optionally scoped to one entity (rtIdScope). Requires StreamDataAdmin.
PeriodicRecomputeOrchestratorHostedService (per-tenant tick)Drains the dirty-windows ledger and recomputes only the windows actually marked dirty — not the whole archive.
Chain propagationinternalAfter 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:

  1. Each rollup table carries a generation column (BIGINT NOT NULL DEFAULT 0) that is part of the primary key. Normal forward aggregation always writes generation = 0.
  2. 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.
  3. A tiny per-rollup side-table (archive_<rtId>__genmap) holds the active generation per range. Flipping that pointer to N+1 is a single-row write — the atomic commit point.
  4. The read path injects generation = CASE WHEN <range> THEN <activeGen> … ELSE 0 END, so every query resolves exactly one consistent generation per window.
  5. 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 shapeavailableFromavailableTo
windowed (rollup, time-range)MIN(window_start) — the start of the earliest stored windowMAX(window_end) — the end of the latest stored window (the exclusive end of the last bucket, not its start)
rawMIN(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:

SurfaceCall
GraphQLstreamData.coverageFor(rtId)
RESTGET /api/v1/streamData/archives/{archiveRtId}/coverage
MCPget_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.

A disabled finer rung can make a coverage end misleading

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:

  1. A rung covers the request when its availableFrom is at or before the requested start (equal counts as covering; the requested end is not considered). A rung without coverage never covers.
  2. The base rung takes part in the filter like every other rung.
  3. 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.
  4. Candidates are the covering rungs; if none covers, the candidates are the rungs with the earliest availableFrom, ties kept.
  5. 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@1 for rows that already carry their source rtId, SaveTimeRangeSeriesInArchive@1 for 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 double mXparser computes: Boolean, Int, Int64, Double, or DateTime (ticks). String and non-scalar types are not supported. A formula / NaN / NULL-operand error stores NULL for that cell; ingest never fails.
  • CK model — CkArchiveColumn carries Name + Formula + ResultType + the engine-managed ComputedState (since System.StreamData 1.5.0) and ComputedVersion / PendingFormula (since 1.6.1 / 1.6.2 respectively). A column with a Formula is computed; one with a Path is 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 to Active atomically. A backfill failure leaves it Failed and 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:

MutationReturnsDescription
activateArchive(rtId)ArchiveTransitionResultProvisions the CrateDB table, transitions to Activated.
disableArchive(rtId)ArchiveTransitionResultActivated → Disabled. Table is preserved.
enableArchive(rtId)ArchiveTransitionResultDisabled → Activated. Re-validates column paths.
retryArchiveActivation(rtId)ArchiveTransitionResultFailed → Activated. Idempotent (CREATE TABLE IF NOT EXISTS).
deleteArchive(rtId)BooleanDrops the CrateDB table and soft-deletes the entity. Rejected if any rollup lists it as a source.
createRollupArchive(input)OctoObjectIdsources[] + aggregations; server-derived TargetCkTypeId + Columns.
createTimeRangeArchive(input)OctoObjectIdOperator-defined target type and columns.
freezeRollupArchive(rtId, until)ArchiveTransitionResultMonotonic FrozenUntil set.
unfreezeRollupArchive(rtId, acceptGaps)ArchiveTransitionResultClears FrozenUntil.
rewindRollupWatermark(rtId, toBucketEnd)ArchiveTransitionResultRe-aggregates from the given boundary.
recomputeArchive(rtId, from, to, rtIdScope?)RecomputeJobInfoTriggers/coalesces an optimistic recompute of [from, to), optionally per entity. Returns the job snapshot.
addComputedColumn(rtId, name, formula, resultType, indexed?)ArchiveTransitionResultAdds a computed column and backfills it (optimistic / atomic).
updateComputedColumnFormula(rtId, name, formula)ArchiveTransitionResultRe-formulates a computed column (versioned backfill + atomic swap). Rejected when another computed column references it.
removeComputedColumn(rtId, name)ArchiveTransitionResultRemoves 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​

MethodPathDescription
POST/enable / /disableTenant-level toggle. /disable answers 409 while any archive is still Activated (AB#4255).
POST/archives/{archiveRtId}/activateSame as GraphQL activateArchive.
POST/archives/{archiveRtId}/disable / /enable / /retryStatus transitions.
DELETE/archives/{archiveRtId}Same as GraphQL deleteArchive.
POST/archives/{rollupRtId}/freeze / /unfreeze / /rewindRollup-only controls.
POST/archives/{rollupRtId}/recomputeOptimistic recompute of [from, to) (rtIdScope optional). Returns the job snapshot.
GET/archives/{archiveRtId}/recompute-jobsRecent recompute jobs (newest first) for debugging.
POST/archives/{archiveRtId}/insertTimeRangeBulk time-range insert.
GET/archives/{archiveRtId}/rollupsActive rollups referencing this archive as a source.
GET/archives/{archiveRtId}/coverageMeasured 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 source rtId, reads the window from two attribute paths, and verifies that every distinct source rtId in 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, is Activated and 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​

  1. Resolve archive (single lookup), status check.
  2. Pre-validate every point in the batch — required-path coverage, type compatibility, path resolution. Stops at the first violation.
  3. On violation → throw with the offending index. No row written.
  4. 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.
  • RawRetentionMs on the Archive base type defines the retention horizon. null = retain forever.
  • The ArchiveRetentionScheduler (IHostedService, hourly) drops partitions older than now() - RawRetentionMs per 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:

  1. All DDL is idempotent (IF EXISTS / IF NOT EXISTS).
  2. Operation order is Crate-first, Mongo-last: Activate runs CREATE TABLE, then flips Status in Mongo. Delete runs DROP TABLE, then soft-deletes the entity. Mongo is the source of truth.
  3. Every data-plane call checks Status == Activated before touching Crate. While step 1 has run but step 2 has not, the API rejects with ArchiveNotActivatedException.
  4. On startup the ArchiveReconciler enumerates Activated archives per tenant and (re-)provisions any missing CrateDB tables. Drift in the other direction (table without entity) is logged but not auto-dropped.

Concurrency​

ScenarioBehavior
Two parallel Activate on the same archiveSerialized via RepositoryDistributedLockService. The second call observes Activated and returns idempotently.
Disable while inserts are in flightIn-flight inserts complete. New inserts after the flip raise ArchiveNotActivatedException.
Delete while inserts are in flightSame as Disable, plus the DROP TABLE may retry once with a short backoff if the table is briefly locked.
Activate fails mid-DDLStatus → Failed. EnsureArchiveCreatedAsync is idempotent (CREATE TABLE IF NOT EXISTS), so retry is safe.
Bulk insert on disabled archiveWhole 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.

ExceptionThrown when
StreamDataExceptionBase — generic StreamData failure.
StreamDataNotEnabledExceptionInstance or tenant flag is false on an operation that requires it.
ArchiveNotFoundExceptionarchiveRtId does not resolve to an Archive entity.
ArchiveNotActivatedExceptionArchive exists but Status ≠ Activated on insert / query.
ArchiveSchemaImmutableExceptionPost-activation update tries to change schema-relevant fields.
ArchivePathInvalidExceptionA CkArchiveColumn.Path cannot be resolved against TargetCkTypeId.
ArchiveColumnTypeUnsupportedExceptionAn attribute type cannot be mapped to a CrateDB column type.
RequiredAttributeMissingExceptionInsert lacks a value for a required: true scalar path. Carries the offending point index.
ArchiveActivationFailedExceptionDDL execution failed during activation. Wraps the underlying SQL error.
InvalidArchiveStateTransitionExceptionLifecycle mutation called against an incompatible status.
RollupSourceInUseExceptiondeleteArchive invoked on an archive that is still listed as a source of a rollup.
RollupSourceMissingExceptionA declared source archiveRtId does not exist or is soft-deleted.
RollupSourceNotActivatedExceptionA declared source archive exists but is not Activated.
RollupSourcesRequiredExceptionA rollup declares no source at all.
DuplicateRollupSourceExceptionThe same source archive is listed more than once.
RollupSourceSpanInvertedExceptionA bounded validity span has ValidFrom >= ValidTo.
RollupSourceSpanOpenEndConflictExceptionMore than one source has an open start, or more than one an open end.
RollupSourceSpanOverlapExceptionTwo validity spans overlap.
RollupSourceSpanNotOnBucketBoundaryExceptionA span boundary is not on the rollup's bucket grid in its reference time zone.
RollupSourceTargetTypeMismatchExceptionA source targets a different CK type than the rollup.
RollupSourcePathMissingExceptionAn aggregation path/function pair does not resolve on one of the sources.
RollupSourceCycleException / RollupCycleExceptionThe declared sources form a cycle, or the rollup lists itself.
RollupSourceDeclarationConflictExceptionThe deprecated SourceArchiveRtId disagrees with a non-empty Sources list.
RollupBucketIntervalExceptionA source's granularity is incompatible with the rollup's bucket size / alignment.

Observability​

Metrics​

MetricTypeTags
streamdata.archive.countGaugetenant, status
streamdata.archive.status_transitionsCountertenant, archive, from, to
streamdata.archive.activation_duration_msHistogramtenant, archive, outcome
streamdata.insert.duration_msHistogramtenant, archive, batch_size_bucket
streamdata.insert.pointsCountertenant, archive
streamdata.insert.required_violationsCountertenant, archive
streamdata.query.duration_msHistogramtenant, archive, query_type
streamdata.query.rows_returnedHistogramtenant, archive, query_type
streamdata.crate.connections.openGauge(instance-wide)
streamdata.crate.retention_partitions_droppedCountertenant, archive
streamdata.crate.rollup_rows_insertedCountertenant, 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​

  1. One archive per concern. Different retention, different query patterns, different consumers → different archives.
  2. Be explicit about required. Required scalars become NOT NULL columns and are validated app-side before insert.
  3. Index sparingly. Indexes cost storage and insert throughput. Default to indexed; opt out for columns that are read but never filtered.
  4. Roll up early. A RollupArchive is far cheaper to query over long horizons than a raw scan.
  5. Pick alignment for the use case. Energy / EDA workflows want calendar alignments (monthly, weekly); engineering telemetry usually wants FixedSize.
  6. 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.
  7. Soft-delete first. Disable an archive to confirm no consumer breaks before deleting.

See Also​