Stream Storage: Append-Only Segments for High-Volume Kinds

Status: Implemented
Purpose: Avoid CAS overhead for high-volume kinds by writing to append-only segment files (multiple records per file). Uses the change journal pattern where only deltas are stored for change_journal_entry. The legacy one-file-per-object (CAS) format is deprecated for high-volume kinds; see HIGH_VOLUME_STORAGE_DEPRECATION.md. References: BYPASS_KIND_STORAGE.md, AUDIT_STREAM_FORMAT.md, INTERNAL_OBJECTS_AS_DEDICATED_WALS.md, HIGH_VOLUME_EVENT_INDEXES.md.

Vocabulary: End-to-end policy-driven maintenance of these streams (compaction, retention, overlays) is object stream stewardship (synonym: stream maintainer pattern). Requirement articulation: REQ-STREAM-001. Template for new environments: scripts/templates/glossary_object_stream_stewardship.yaml.

Implementation hook: After retention, the scheduler maintenance runner invokes storage.PostRetentionStreamStewardship (pkg/storage/stream_stewardship.go) β€” registry compaction, orphaned segment GC, runtime-delta backfill and overlay GC. Steward wake is per stream-backed kind (see STREAM_KIND_STEWARDSHIP.md). Same entry point is safe to call from tests on an empty tree.


1. Enablement

  • Env: ZQK_STREAM_STORAGE_ENABLED=true (or 1) enables stream storage for supported kinds.
  • Kinds: audit_event, change_journal_entry, all high-volume metrics (audit_aggregation_metric, base_metric, command_metric, file_lock_metric, code_quality_metric, scheduler_health_metric), mcp_session, zqk_session, verification_matrix. See high_volume_kinds.yaml and stream_config.go.
  • When enabled, Create for those kinds writes to stream segments instead of CAS (no hash file, no hash registry, no CAS index update). Retention continues to use the same index contract (Count, OldestIDs, IDsOlderThan) via the high-volume event cache, which is updated on stream append.

2. Segment layout

All stream-backed kinds are organized under .zqk/streams/<kind>/ with UTC daily segment rolling:

Kind Identifier Segment Directory Canonical Filename Pattern Description & Notes
audit_event .zqk/streams/audit_event/ <yyyy-mm-dd>_stream.json High-volume append-only audit stream. Legacy audit_stream_YYYY-MM-DD.jsonl files are automatically migrated to canonical format.
change_journal_entry .zqk/streams/change_journal_entry/ <yyyy-mm-dd>_stream.json Delta-only change journal records.
*_metric (High-volume metrics) .zqk/streams/<kind>/ <yyyy-mm-dd>_stream.json Covers audit_aggregation_metric, base_metric, command_metric, file_lock_metric, code_quality_metric, and scheduler_health_metric.
verification_matrix .zqk/streams/verification_matrix/ <yyyy-mm-dd>_stream.json Verification matrix snapshots; active updates resolve via stream_current (stream_current_state.go).
mcp_session / zqk_session .zqk/streams/<kind>/ <yyyy-mm-dd>_stream.json Session lifecycle event journals.
  • Rolling Window: Daily (UTC). One primary file per calendar day; next day rolls to a new file. When concurrent worker processes write to the same daily stream, process IDs are disambiguated as <yyyy-mm-dd>_pid<pid>_stream.json.
  • Format: JSONL (one JSON object per line, UTF-8 encoded, \n line terminator).
  • Concurrency: Append-only with per-kind thread synchronization mutex.

3. Change journal pattern (delta-only)

change_journal_entry: Only delta fields are written (no full-state bloat):

  • Stored: id, object_ref, change_type, diff_summary, changed_paths, created_at, created_by, title, and optionally previous_state (for rollback).
  • Omitted from stream: Full object snapshot and other fields not needed for aggregation or reconstruction.

Metrics (*_metric): Minimal delta-style fields: id, kind, status, created_at, created_by, plus kind-specific (e.g. window_start, window_end, aggregated_entry_count for audit_aggregation_metric). Retention uses the same index contract (Count, OldestIDs, status for protect_statuses); segment files stay small.

See INTERNAL_OBJECTS_AS_DEDICATED_WALS.md and change_journal_compaction.md for the same delta/compaction philosophy.


4. Index and retention

  • Per-record location: On append, the write path registers id β†’ segmentPath::offset in the stream location registry (in-process) and adds an entry to the high-volume event cache (id, kind, created_at, file_path = segmentPath::offset).
  • Same contract: Count, OldestIDs, and IDsOlderThan are satisfied by the high-volume cache; retention and aggregation logic are unchanged.
  • Delete: Stream-backed objects are soft-deleted: removed from the stream registry and from the cache; the segment file is not rewritten (record remains on disk but is no longer visible to List/Get/retention).

5. Implementation

  • Append: pkg/storage/stream_storage.go – AppendToStream(projectRoot, kind, id, obj, createdAt) β†’ (segmentPath, offset). recordForStream reduces change_journal_entry to delta fields and metrics to minimal delta-style fields.
  • Config: pkg/storage/stream_config.go – StreamStorageEnabledForKind(kind).
  • Path resolution (no hardcoding): All paths use the path alias cache and scheme prefixes (prefix:, abs:, web:). See PATH_ALIAS_RESOLUTION.md. Segment dirs resolve as prefix:streams/<kind>. Call sites use paths.ResolvePath or GetStreamSegmentDir. The cache is built during pre-warm Tier 0. Architectural policy and requirement: see process objects (path resolution from cache). Do not hardcode paths; use the resolver.
  • Create path: object_storage_file_create.go – when stream enabled for kind, writeObjectToStream appends, sets stream location, and updates high-volume cache; no CAS or hash registry.
  • Object path resolution: getObjectFilePath returns segmentPath::offset when the object is in the stream registry.
  • Delete: object_storage_file_delete.go – if path is segmentPath::offset, soft delete (remove from registry and cache only).

6. Alignment with timeseries prototype (delta storage)

The codebase has a timeseries prototype for numeric, append-only data that stores only deltas (base + delta encoding, chunked files):

  • pkg/metrics/timeseries.go: TimeSeriesWriter / TimeSeriesReader – chunked files per time window, delta time (varint seconds from chunk base) and delta value (zigzag varint). Used for object_volume (object-count-report) and health-check monitors. See TODO(PRI-218 / metrics-timeseries) for integration checklist.

Alignment:

Use case Storage shape Delta pattern
Event log (audit_event, change_journal_entry) Stream segments (JSONL per day) change_journal: only delta fields (object_ref, change_type, diff_summary, changed_paths, …). audit_event: full event per line (could add dictionary compression later).
Numeric time series (object_volume, gauges) Chunked binary (timeseries.go) Base + delta time and delta value (varint/zigzag).

Both patterns avoid full-state bloat: event streams store minimal per-record fields (delta-only for change_journal); time series store deltas from a chunk base. For metrics objects (audit_aggregation_metric, *_metric), stream storage uses the same segment-per-day layout and minimal per-record fields (id, kind, status, created_at, and only what retention/aggregation need) so storage stays optimized.


7. Stream-only write policy (no YAML for stream-backed kinds)

When stream storage is enabled for a kind, that kind must not be written to .zqk/process/ YAML (CAS). All writes must go through the stream path so that:

  • No new files are created under .zqk/process/audit/, .zqk/process/scheduler_jobs/, etc., for those kinds.
  • Object count and retention use the high-volume cache and stream registry only; disk YAML count for those dirs does not grow.

Required behavior:

  1. Create / update / aggregated flush: Use fileStorage.Create, fileStorage.Update, or writeObjectToStorage so the code path chooses stream when StreamStorageEnabledForKind(kind) is true. Do not use WriteSystemObjectAndRegisterHash for stream-backed kinds when stream is enabledβ€”that would create YAML under .zqk/process. See audit_events.go (WriteSystemObjectAndRegisterHash comment) and audit_event_buffer_flush.go (aggregated event flush).
  2. Audit event buffer flush: When stream is enabled for audit_event, the buffer requires fileStorage and uses only writeObjectToStorage (stream). If fileStorage is nil, flush returns an error instead of falling back to CAS/YAML.
  3. Write-behind replay: WAL replay for stream-backed kinds uses minimal records (no payload); apply skips when data is empty and kind is stream-enabled. Replay with payload still goes through writeObjectToStorage β†’ stream.

Stream-enabled kinds (see stream_config.go): audit_event, change_journal_entry, audit_aggregation_metric, base_metric, command_metric, file_lock_metric, code_quality_metric, scheduler_health_metric, mcp_session, scheduler_job, zqk_session.


8. Stream registry and deleted-ID state files (growth and cleanup)

Location: .zqk/state/
Files (per kind):

  • stream_registry_<kind>.jsonl β€” Append-only. Each line: {"id":"...","loc":"segmentPath::offset"}. Every created stream-backed object gets one line; lines are never removed on delete.
  • stream_deleted_<kind>.jsonl β€” Append-only. One object ID per line. Every soft-deleted stream-backed object gets one line; lines are never removed.

Why both: On delete we only remove the ID from the in-memory registry; we do not rewrite the registry file. So the registry file still contains idβ†’loc for deleted events. The deleted file is the source of truth for β€œthis ID was deleted”: when loading the snapshot we build locations from the registry and deleted from the deleted file, and treat an ID as live only if it’s in locations and not in deleted. So both files are required for correct List/Get/Count.

Current behavior: Both files grow indefinitely. There is no retention, compaction, or truncation. For high-volume kinds (e.g. audit_event) with aggressive retention, the deleted file can reach hundreds of thousands of lines (e.g. 400k+). Impact:

  • Disk: State directory size grows without bound.
  • Load cost: loadStreamRegistrySnapshot reads both files in full and builds two maps. Cache is invalidated after every batch of deletes, so large deleted files mean repeated full reads and large in-memory sets.

Cleanup plan (to implement): Add a compaction path (scheduler job or maintenance command) that:

  1. Reads the full registry and full deleted set for the kind.
  2. Writes a new registry file containing only lines for IDs that are not in the deleted set (i.e. live IDs only).
  3. Replaces the old registry with the new one (atomic rename).
  4. Truncates or rewrites the deleted file to empty (or to a minimal sentinel), since all those IDs are no longer in the registry and no longer need to be tracked.

Routine compaction: Run zqk system compact-stream-state --kind audit_event (or other stream-backed kind) when the deleted count is high. See WAL_AND_STATE_FILES.md for all WAL and state file naming and cleanup.


9. References

  • Glossary (logical data paths): Data stream summary, Data cell (coordinator membrane around stream/CAS nucleus). Requirements: REQ-DATASTREAM-001, REQ-DATACELL-001.
  • DATA_STORAGE_PRODUCTION_ROADMAP.md β€” Optimize first, then timestamp compression and timeseries wiring
  • BYPASS_KIND_STORAGE.md – Motivation and aggregate storage
  • AUDIT_STREAM_FORMAT.md – Audit stream file format
  • INTERNAL_OBJECTS_AS_DEDICATED_WALS.md – WAL pattern for internal kinds
  • HIGH_VOLUME_EVENT_INDEXES.md – Index contract (Count, OldestIDs, IDsOlderThan)
  • pkg/metrics/timeseries.go – Timeseries prototype (base+delta, chunked)
  • pkg/storage/stream_storage.go, stream_config.go, high_volume_event_cache.go
  • Β§7 above: stream-only write policy (no YAML for stream-backed kinds)
  • WAL_AND_STATE_FILES.md β€” WAL and state file naming, lifecycle_events purpose, routine compaction