Skip to content

About

Kubernetes-native data consistency SLO operator for derived data correctness

Resources

Stars

0 stars

Watchers

0 watching

Forks

Repository files navigation

KubeDataGuard

One-Line Idea

KubeDataGuard is an evidence-first data consistency SLO engine with an optional Kubernetes operator for scheduling checks, publishing compact status, and triggering safe reconciliation workflows.

Current Maturity

This repository is a research prototype vertical-slice MVP, not a production-ready general-purpose data platform yet.

What is implemented and verified:

  • Postgres -> Redpanda/Kafka -> OpenSearch demo pipeline.
  • Second derived-view path: Postgres -> Redpanda/Kafka -> Redis cache freshness.
  • Analytics derived-view path: Postgres -> ClickHouse paid-order aggregate checks.
  • Hardcoded commerce invariants for existence, aggregate consistency, and bounded freshness.
  • Prefix-bucket checksum/Merkle-style narrowing for large Postgres -> OpenSearch comparisons.
  • A first generic query invariant for Postgres source queries and OpenSearch JSON target queries.
  • Keyset-paginated Postgres source scans for query invariants, with source scan evidence in reports.
  • Persisted source-scan checkpoints for bounded query scans, stored in the same local or S3-compatible report store.
  • OpenSearch target query pagination with search_after; size is treated as page size, not a correctness cap.
  • First CDC frontier proof for the commerce path: Postgres WAL LSN, outbox publication state, Kafka partition offsets, and OpenSearch applied-offset evidence are recorded in the observation window.
  • A Go/controller-runtime operator that schedules checker and repair Jobs.
  • A Kubernetes-native kind demo stack for Postgres, Redpanda, OpenSearch, Redis, ClickHouse, checker Jobs, and the operator.
  • Production-shaped demo StatefulSets with PVCs, probes, and resource bounds for Postgres, Redpanda, OpenSearch, Redis, and ClickHouse.
  • Compact Kubernetes status handoff through ConfigMaps.
  • Full JSON/Markdown reports written to the worker report store and referenced from status.
  • Optional S3/MinIO-compatible durable report publication with s3://... status refs.
  • Kubernetes Secret-backed connection env resolution through DataSource and DerivedView resources.
  • DataSource, DerivedView, and Secret rotation watch fan-out to dependent Invariant reconciles.
  • Direct source-of-truth reindex repair for explicit local demos.
  • Safer owner-executed repair modes: reconciliation JSONL, Kafka replay requests, webhook dispatch, Redis invalidation requests, and ClickHouse backfill requests.
  • Prometheus textfile metrics and OTEL-shaped JSONL span artifacts for check/repair reports.
  • Compact report artifacts, local retention cleanup, and optional S3 server-side encryption metadata.
  • Operator safeguards that avoid status churn for repeated healthy scheduled checks.

What is intentionally not solved yet:

  • Arbitrary databases and target stores beyond the current Postgres/OpenSearch, Redis-cache, and ClickHouse analytics paths.
  • Full DBLog-style logical decoding, exact snapshot-plus-log recovery, and crash-exact scan recovery.
  • Managed object-store lifecycle policies outside KubeDataGuard's own upload metadata.
  • Fully managed production operators for databases; the checked-in stack is production-shaped, not a replacement for CloudNativePG, Redpanda Operator, or OpenSearch Operator.

Architectural Boundary

The first-principles rule is:

Kubernetes is the control plane, not the data plane.

KubeDataGuard should use Kubernetes to declare invariants, schedule workers, expose compact semantic status, and connect to normal platform alerting. It should not use etcd, CRD status, or ConfigMaps as a high-volume row-level data store.

The long-term shape is:

Invariant CRD / CLI config
  -> checker worker
  -> durable report store such as S3, GCS, or MinIO
  -> compact health metric/status only
  -> owner-approved reconciliation request

That boundary matters because a real drift report can contain thousands or millions of counterexamples. Kubernetes should hold the pointer and the phase, not the blob. The checked-in operator keeps ConfigMaps compact, avoids rewriting Invariant.status when only telemetry such as checkID, reportRef, or checked row count changes, and can publish full evidence reports to an S3-compatible store.

Evidence Storage

Local files remain the default:

REPORT_STORE=local
REPORT_DIR=/workspace/reports

For durable evidence, configure an S3-compatible sink:

REPORT_STORE=s3
REPORT_BUCKET=kubedataguard-reports
REPORT_PREFIX=prod/orders
REPORT_S3_ENDPOINT_URL=http://minio:9000   # optional, for MinIO/S3-compatible stores
AWS_REGION=us-east-1

The worker always writes local JSON/Markdown artifacts first, then uploads the report, Markdown summary, and latest.json when REPORT_STORE=s3. Invariant.status.reportRef becomes an s3://bucket/key pointer while the full counterexample evidence stays outside Kubernetes. Query scan checkpoints use the same backend under scan-checkpoints/, so long scans can be bounded per run without losing the last processed key.

Research-Informed Novelty Wedge

The ecosystem already has excellent infrastructure operators and data-quality tools. CloudNativePG manages PostgreSQL lifecycle and failover. Strimzi manages Kafka. KubeDB and KubeBlocks manage many databases and day-2 operations. Great Expectations and Soda validate data quality. OpenLineage records lineage metadata.

KubeDataGuard should not compete with any of those directly.

Its sharper contribution is:

Evidence-backed reconciliation of cross-system derived-data correctness.

That means the project should become the missing layer between infrastructure health and data quality:

  • Infrastructure operators ask: is the database, stream, or index healthy?
  • Data-quality tools ask: does a dataset satisfy expectations?
  • KubeDataGuard asks: does this derived view still correctly represent that source of truth within a declared SLO?

The most novel version of this project is not "a checker." It is a data-correctness control plane with:

  • DataSource, DerivedView, Invariant, and RepairPolicy CRDs
  • lineage-bound invariants that know how data is supposed to move
  • Invariant.status as compact Kubernetes truth when the operator wrapper is used
  • durable drift reports as evidence
  • repair provenance showing exactly what was requested and why
  • checker failure states that are distinct from data-drift states
  • OpenTelemetry traces for check and repair cycles
  • optional OpenLineage emission so the data ecosystem can see the same lineage graph

Demo Story

The first demo should run in a few commands:

make demo-local # no-Docker proof: drift -> report -> repair -> clean verification
make up       # starts Postgres, Redpanda, Redpanda Console, OpenSearch, and Redis
make seed     # resets the demo and inserts 50 orders, most marked paid
make drift    # commits some paid events without indexing them
make drift-cache # commits some paid events without updating Redis
make check    # reports missing paid orders with order IDs and source details
make check-checksum # compares source/target bucket fingerprints, then drills into bad buckets
make check-redis-freshness # checks the Postgres -> Redis cache freshness invariant
make repair   # reindexes missing/stale orders from Postgres and verifies
make check    # produces a clean report
make demo-freshness-drift # proves bounded-freshness drift with source/target timestamps
make up-analytics && make demo-clickhouse-drift # proves analytics aggregate drift in ClickHouse
make k8s-demo # runs the Kubernetes-native kind stack and operator proof

That is the project in miniature: a derived view becomes untrustworthy, KubeDataGuard proves exactly how, then repairs and verifies it.

The demo-local target remains useful as a fast proof path without external services. The full Compose path has also been runtime-verified with Docker CLI + Colima, Postgres, Redpanda, and OpenSearch.

First Principles

Modern applications rarely keep data in one place. A single user action can write to Postgres, emit an event to Kafka, update Redis, index a document into Elasticsearch, and update analytics tables in ClickHouse. Each copy exists for a performance or workflow reason, but each copy can become wrong.

The fundamental problem:

If the same business fact exists in multiple systems, how do we know those systems still agree?

Traditional monitoring says:

Is Postgres up?
Is Kafka up?
Is Elasticsearch responding?
Is latency acceptable?

KubeDataGuard asks a deeper question:

Is the data correct enough for the business promise we made?

The project is built around four ideas:

  1. Source data

    • The first durable record of a business fact.
    • Example: orders table in Postgres.
  2. Derived data

    • A transformed, copied, cached, indexed, or aggregated view.
    • Example: orders index in Elasticsearch.
  3. Invariant

    • A rule that must remain true.
    • Example: every paid order in Postgres must appear in the search index within 60 seconds.
  4. Reconciliation

    • A controlled action that asks the owning pipeline to make a derived view trustworthy again.
    • Example: publish a reconcile event, call an owner service webhook, replay Kafka events, invalidate Redis keys, run a backfill job, or, in a local demo only, direct reindex from the source.

DDIA Concepts

This project is strongly connected to:

  • Replication and replication lag
  • Derived data and denormalized views
  • Stream processing
  • Batch backfills
  • Encoding and evolution
  • Transactions and the absence of distributed atomicity
  • Fault tolerance
  • Observability and operability

The key DDIA insight is that many useful data systems are eventually consistent, but "eventually" is not an engineering plan. KubeDataGuard turns eventual consistency into explicit, measurable SLOs.

Academic Foundation

The research-backed version of KubeDataGuard should be more formal than a table-diff tool.

The consistency-model papers push the project toward named guarantees instead of generic "drift." An invariant should declare which claim it protects:

  • Existence: every qualifying source record appears in the derived view.
  • Equality: selected source fields and target fields agree.
  • Aggregate: counts, sums, or grouped totals agree inside a tolerance.
  • Freshness: the derived output is no older than the allowed lag window.
  • Client-centric consistency: read-your-writes and monotonic-read claims hold for a traced client/session.
  • Causal consistency: if fact B depends on fact A, the derived system must not expose B without A.

DBLog is the most important systems paper for the MVP. It shows why snapshot-plus-log CDC needs watermarks, chunk boundaries, and resumable scans. The commerce path now records a first frontier chain: source WAL LSN, source watermark, outbox publication counts, per-partition Kafka offsets, target applied-offset evidence, and target read time. That is not full logical-decoding CDC yet, but it stops the report from pretending that a moving source and moving target were compared without boundaries.

The polyglot-persistence papers sharpen the novelty claim. The problem is not only that Postgres, Kafka, OpenSearch, Redis, and ClickHouse exist. The problem is that business truth leaks across them, and each store has different query semantics, failure modes, indexes, and freshness behavior. KubeDataGuard's CRDs are valuable only if they can describe those heterogeneous links without hiding the evidence.

Elle changes the reporting philosophy. A good report should not say only:

Invariant failed: 12 missing records

It should give a counterexample:

order-123 was committed in Postgres at source version 981,
published to topic order-events at offset 4412,
but was absent from OpenSearch after the 60 second lag window.

GITER and the Kubernetes operator model reinforce the final control-plane shape:

spec = declared consistency promise
status = compact observed truth
report = durable evidence
repair = controlled reconciliation action

DCaaS is the conceptual ancestor: data correctness policy should live outside application code. KubeDataGuard updates that idea for Kubernetes, CDC pipelines, polyglot stores, and SLO-driven repair.

Kubernetes Angle

Kubernetes already has a powerful reconciliation model:

desired state -> controller observes actual state -> controller acts -> status is updated

KubeDataGuard applies this model to data correctness:

desired consistency -> operator observes actual data -> operator detects drift -> operator repairs or reports

Instead of only declaring deployments and services, teams declare data contracts as Kubernetes resources.

Implementation Language Direction

Use a hybrid architecture:

Go/controller-runtime operator
  |
  |-- watches DataSource, DerivedView, Invariant, RepairPolicy
  |-- schedules checker/repair Jobs
  |-- writes compact Kubernetes status
  |
Python checker/repair workers
  |
  |-- connect to Postgres, Kafka/Redpanda, OpenSearch, Redis, ClickHouse
  |-- run invariant logic
  |-- emit verbose evidence reports
  |-- emit compact status payloads for the operator

Python is the right language for the local MVP and data-system integrations because the iteration loop is fast, client libraries are good, and invariant logic is easier to evolve. Go is the right language for the long-term Kubernetes operator because controller-runtime, CRD status handling, leader election, watches, work queues, and Kubernetes API machinery are strongest there.

The important boundary is the report contract. Python workers now emit:

  • verbose JSON/Markdown evidence reports
  • compact kubernetes_status payloads shaped like Invariant.status

That means the Go controller can remain thin: reconcile resources, run workers, store report references, and update status.

Current closed-loop operator path:

  • cmd/dataguard-operator: controller-runtime manager entrypoint
  • internal/controller: unstructured Invariant reconciler
  • watches dataguard.io/v1alpha1 Invariant, DataSource, and DerivedView resources
  • maps DerivedView changes to referencing Invariants and maps DataSource changes through dependent DerivedViews
  • supports dataguard.io/checker-mode=job
  • creates checker Jobs per Invariant generation or scheduled check interval
  • runs the Python checker against Postgres, Redpanda/Kafka, and OpenSearch
  • stores compact status.json and repair input in a report ConfigMap
  • writes full JSON/Markdown reports to the worker report store and keeps only the report reference in status
  • if drift is detected, finds an explicitly allowed RepairPolicy
  • never auto-runs direct reindex-records unless the policy has approvalRequired: false and the unsafe opt-in annotation dataguard.io/allow-unsafe-direct-reindex: "true"
  • repair Jobs can run direct reindex for controlled demos or emit reconciliation requests for safer owner-executed repair
  • copies compact checker or repair status.json into Invariant.status only when semantic health changes
  • still keeps synthetic mode available for controller smoke tests

This proves the Kubernetes-native control loop without making Kubernetes the row-level data plane: the operator reconciles a declared data SLO, a Job performs the data check, the full evidence lives outside CRD status, and the CRD status reflects only the compact verified consistency state.

Core Custom Resources

DataSource

Describes a system that stores data.

apiVersion: dataguard.io/v1alpha1
kind: DataSource
metadata:
  name: orders-postgres
spec:
  type: postgres
  connectionSecret: orders-postgres-secret
  connectionSecretKey: dsn
  database: commerce

DerivedView

Describes how one data store is derived from another.

apiVersion: dataguard.io/v1alpha1
kind: DerivedView
metadata:
  name: orders-search-index
spec:
  sourceRef: orders-postgres
  target:
    type: opensearch
    connectionSecret: search-secret
    connectionSecretKey: url
    index: orders
  pipeline:
    type: kafka
    connectionSecret: kafka-secret
    bootstrapServersKey: bootstrapServers
    topic: order-events

Invariant

Describes the rule that must hold.

apiVersion: dataguard.io/v1alpha1
kind: Invariant
metadata:
  name: paid-orders-indexed
spec:
  derivedViewRef: orders-search-index
  type: existence
  checkIntervalSeconds: 300
  maxLagSeconds: 60
  sourceQuery: |
    select id from orders where status = 'paid'
  targetQuery: |
    select id from orders_search where status = 'paid'

The first executable generic query invariant is intentionally narrower:

apiVersion: dataguard.io/v1alpha1
kind: Invariant
metadata:
  name: paid-orders-query-check
spec:
  derivedViewRef: orders-search-index
  type: query
  checkIntervalSeconds: 300
  maxLagSeconds: 60
  keyField: id
  sourceScanPageSize: 1000
  compareFields:
    - status
    - amount_cents
    - currency
    - version
  sourceQuery: |
    select id, status, amount_cents, currency, version, updated_at
    from orders
    where status = 'paid'
      and updated_at <= now() - interval '60 seconds'
    order by updated_at asc, id asc
  targetQuery: |
    {
      "query": {
        "bool": {
          "filter": [
            {"term": {"status": "paid"}},
            {"range": {"updated_at": {"lte": "now-60s"}}}
          ]
        }
      }
    }

For intentionally bounded long scans, add sourceCheckpointId and sourceMaxPages; those runs are marked partial until a complete snapshot/checkpoint model exists.

RepairPolicy

Describes what may happen when an invariant fails.

apiVersion: dataguard.io/v1alpha1
kind: RepairPolicy
metadata:
  name: reconcile-missing-paid-orders
spec:
  invariantRef: paid-orders-indexed
  approvalRequired: true
  actions:
    - type: emit-reconcile-events
      batchSize: 500

Direct automatic reindex-records is intentionally fenced off. For a controlled demo or explicitly accepted risk, a policy must set both:

metadata:
  annotations:
    dataguard.io/allow-unsafe-direct-reindex: "true"
spec:
  approvalRequired: false

Reconciliation Loop

The mature operator loop:

1. Watch DataSource, DerivedView, Invariant, and RepairPolicy objects.
2. Resolve connection env vars from `DataSource` and `DerivedView` Secret references.
3. Run source and target checks.
4. Compare results.
5. Classify drift:
   - missing data
   - stale data
   - duplicate data
   - aggregate mismatch
   - schema incompatibility
6. Update Invariant status.
7. Emit metrics and traces.
8. Trigger repair if policy allows it.
9. Verify repair.
10. Record the outcome.

The implemented checker loop today is:

Invariant CRD
  |
  |-- Go operator sees dataguard.io/checker-mode=job
  |-- computes checkID from generation or checkIntervalSeconds time slot
  |-- creates dataguard-check-<invariant>-<checkID> Job
  |-- Python checker connects to Postgres, Redpanda/Kafka, and OpenSearch
  |-- Python checker writes full report to REPORT_DIR and optionally S3/MinIO
  |-- Python checker writes compact status.json and repair-input.json into a ConfigMap
  |-- Go operator copies status.json into Invariant.status
  |-- Go operator requeues scheduled invariants for the next check interval

This path now has a Kubernetes-native demo stack under k8s/demo-stack.yaml, so kind can run Postgres, Redpanda, OpenSearch, checker Jobs, and the operator without reaching through host.docker.internal.

Invariant Evidence Model

An invariant is not just a query pair. The durable model should be:

Invariant =
  scope
  + consistency claim
  + observation window
  + allowed lag
  + evidence
  + repair policy

Evidence should be concrete enough for a human or controller to reproduce the finding:

  • source object ids and versions
  • target object ids and versions
  • source transaction timestamp or LSN when available
  • stream topic, partition, and offset when available
  • source snapshot window or CDC watermark
  • observed lag in seconds
  • comparison outcome: missing, stale, duplicate, aggregate mismatch, or unknown
  • checker result: complete, partial, failed, or timed out
  • repair provenance: action, input ids, output ids, and verification result

Kubernetes status should stay compact:

Healthy
DriftDetected
CheckFailed
Repairing
RepairFailed
Unknown

The full JSON/Markdown report should carry the counterexamples and replay boundaries.

MVP

Build the smallest convincing version:

Postgres -> Kafka/Redpanda -> Elasticsearch/OpenSearch

Use case:

Every paid order in Postgres must exist in the search index within 60 seconds.

Components:

  • Demo API that writes orders to Postgres and emits order events
  • Indexer worker that consumes events and writes Elasticsearch/OpenSearch documents
  • KubeDataGuard checker that compares Postgres and the index
  • Repair worker that reindexes missing records
  • Docker Compose for local development
  • Kubernetes CRDs after the Compose demo works

Implemented Local MVP

This folder now contains a runnable local MVP skeleton.

Files:

  • docker-compose.yml: Postgres, Redpanda, Redpanda Console, OpenSearch, and the dataguard CLI container
  • Dockerfile: Python CLI image
  • requirements.txt: Postgres, Kafka, and OpenSearch client libraries
  • src/dataguard: local MVP implementation
  • src/dataguard/local_demo.py: no-Docker proof of drift detection, aggregate mismatch, repair, and clean verification
  • tests: pure invariant unit tests
  • docs/ARCHITECTURE.md: deeper system design
  • docs/RUNBOOK.md: commands for running the demo
  • docs/KNOWLEDGE_BASE.md: official docs consulted
  • k8s/crds: first Kubernetes CRD sketches
  • examples/commerce-consistency.yaml: example declarative consistency policy
  • cmd/dataguard-operator: Go/controller-runtime operator entrypoint
  • internal/controller: synthetic and job-backed Invariant reconciliation
  • Dockerfile.operator: Go operator image
  • k8s/operator.yaml: in-cluster operator Deployment/RBAC sketch

Implemented commands:

python -m dataguard.cli init
python -m dataguard.cli generate
python -m dataguard.cli index
python -m dataguard.cli inject-freshness-drift
python -m dataguard.cli check --invariant existence
python -m dataguard.cli check --invariant aggregate
python -m dataguard.cli check --invariant freshness
python -m dataguard.cli check --invariant query --source-query "..." --target-query "..."
python -m dataguard.cli check --invariant query --source-checkpoint-id orders-query --source-max-pages 10 --source-query "..." --target-query "..."
python -m dataguard.cli check-job --invariant existence
python -m dataguard.cli repair
python -m dataguard.cli repair-job
python -m dataguard.cli demo-local

With Docker Compose:

make demo-local
make up
make seed
make drift
make drift-freshness
make check
make check-aggregate
make check-freshness
make repair
make repair-freshness
make check

Verified Compose behavior on this machine:

make demo-drift
  generated 50 orders
  indexed 42 events
  deliberately skipped 8 paid orders
  existence invariant: DriftDetected, drift_count=8
  aggregate invariant: DriftDetected, source count=41, target count=33

make demo-repair
  repaired 8 missing orders
  post-repair existence invariant: Healthy
  post-repair aggregate invariant: Healthy

Verified Kubernetes checker behavior on this machine:

kind + operator + checker Jobs + data systems

generation 2 after deliberate drift:
  paid-orders-indexed: DriftDetected, drift_count=8
  paid-orders-aggregate: DriftDetected, drift_count=2

generation 3 after repair:
  paid-orders-indexed: Healthy, drift_count=0
  paid-orders-aggregate: Healthy, drift_count=0

generation 4 after restoring the 60-second SLO spec:
  paid-orders-indexed: Healthy, drift_count=0
  paid-orders-aggregate: Healthy, drift_count=0

Verified Kubernetes repair behavior on this machine:

generation 5 after deliberate drift:
  checker Job found 8 missing paid orders
  operator found RepairPolicy reindex-missing-paid-orders
  operator created dataguard-repair-paid-orders-indexed-g5
  repair Job reindexed 8 missing OpenSearch documents from Postgres
  repair Job verified paid-orders-indexed Healthy

fresh aggregate generation after repair:
  paid-orders-aggregate: Healthy, drift_count=0

restored 60-second SLO generations:
  both invariants: Healthy, drift_count=0

Verified failure taxonomy behavior on this machine:

checker failure:
  unreachable Postgres endpoint -> checker Job failed -> Invariant phase CheckFailed
  restored endpoint -> next generation Healthy

repair verification failure:
  aggregate drift + temporary reindex RepairPolicy -> repair Job completed
  verification still found aggregate mismatch -> Invariant phase RepairFailed
  policy removed + data repaired -> next generation Healthy

Verified bounded freshness behavior on this machine:

paid-orders-freshness:
  guarantee: boundedFreshness
  checkedRecords: 41
  phase: Healthy
  sourceLSN: present
  sourceWatermark: present
  streamOffsetStart: 0
  streamOffsetEnd: 250

Verified freshness drift controls in the local path:

make demo-freshness-drift
  indexed paid orders first
  injected 5 target documents whose indexed_at was older than source updated_at
  freshness invariant: DriftDetected, drift_count=5

make repair-freshness
  reindexed the freshness violation candidates from Postgres
  post-repair freshness invariant: Healthy
  report preserved historical freshness SLO breach evidence separately

Verified freshness drift repair through the kind operator:

paid-orders-freshness generation 4:
  checker Job: DriftDetected, drift_count=5
  RepairPolicy: reindex-stale-paid-orders
  repair Job: Healthy, drift_count=0, sloBreachCount=5

paid-orders-freshness generation 5 after restoring maxLagSeconds=60:
  Healthy, checkedRecords=41, drift_count=0, sloBreachCount=5

Compose reports are written to:

${HOME}/.kubedataguard/reports

Current no-Docker proof:

report_type: kubedataguard-local-proof
before existence status: DriftDetected
before aggregate status: DriftDetected
repair action: reindex-records-from-source
after existence status: Healthy
after aggregate status: Healthy

Current report shape:

  • status: compact control-plane phase such as Healthy or DriftDetected
  • guarantee: the semantic claim being checked
  • observation_window: check time, target read time, max lag, eligible source boundary, source LSN, stream topic, stream offset range, and CDC frontier evidence when available
  • observation_window.cdc_frontier: Postgres WAL/outbox, Kafka offset, and OpenSearch applied-offset proof metadata for deciding whether the checked window was bounded, partial, or unavailable
  • observation_window.source_scan: keyset scan mode, key field, page size, page count, row count, first key, last key, resume key, query hash, and checkpoint reference when enabled
  • kubernetes_status: compact status payload used by the job-backed operator
  • checkID: generation or scheduled interval identifier used to make repeated checks idempotent
  • counterexamples: compact evidence for missing, stale, or aggregate-mismatch violations
  • missing, stale, aggregate_mismatches, and freshness_violations: detailed current drift classes
  • freshness_breaches: historical bounded-freshness SLO misses that are preserved as evidence but do not make the current state unrecoverably unhealthy

Failure Scenarios To Demonstrate

  • Indexer is down for two minutes
  • Kafka consumer commits an offset but fails before indexing
  • Elasticsearch rejects a document because of a mapping/schema change
  • Duplicate event causes duplicate derived rows
  • Backfill misses a date partition
  • Redis cache contains stale data after source update

Failure Taxonomy

The system distinguishes data drift from checker/repair failure.

Data and pipeline failures:

  • Missing derived record
  • Stale derived record
  • Duplicate derived record
  • Aggregate mismatch
  • Replication lag beyond SLO
  • Schema or mapping rejection

Checker failures:

  • Source unavailable
  • Target unavailable
  • Checker crashes mid-scan
  • Partial scan
  • Slow scan exceeds check interval
  • Credentials or Kubernetes Secret missing

Repair failures:

  • Repair crashes midway
  • Repair writes duplicates
  • Repair uses stale source data
  • Repair succeeds mechanically but verification still fails
  • Repair loops without reducing drift

Expected behavior:

  • Never mark an invariant healthy after an incomplete check.
  • Separate "checker failed" from "data drift detected."
  • Make repairs idempotent.
  • Verify after every repair.
  • Store compact status in Kubernetes and verbose reports externally.

Milestones

Milestone 1: Runtime-Verified Local Demo

  • Implemented no-Docker proof path with deterministic local source and derived data
  • Runtime-verified Docker Compose with Postgres, Redpanda, OpenSearch, and the Python checker/repair CLI
  • Generate demo orders: verified
  • Detect missing indexed orders: verified
  • Detect aggregate count/revenue drift: verified
  • Print JSON/Markdown drift reports: verified
  • Repair and verify the invariant: verified

Milestone 2: Report and Status Shape

  • Keep JSON/Markdown reports for local CLI
  • Keep ConfigMap handoff compact in Kubernetes: status.json plus repair input, not full report blobs
  • Define Invariant.status fields for Kubernetes:
    • healthy
    • phase
    • guarantee
    • checkStatus
    • driftCount
    • lastCheckedAt
    • observationWindow
    • counterexampleCount
    • reportRef
    • observedGeneration
    • checkID
    • checkIntervalSeconds
    • nextCheckAfter

Milestone 3: More Invariant Types

  • Implemented aggregate checks: paid order count and revenue total
  • Implemented bounded freshness checks using source updated_at, target indexed_at, Postgres LSN, Kafka offset evidence, and target applied-offset evidence
  • Implemented prefix-bucket checksum narrowing for large source/target comparisons
  • Implemented freshness drift injection and freshness repair verification
  • Implemented first generic Postgres/OpenSearch query invariant path
  • Implemented Postgres -> Redis cache freshness as the second source/derived pair
  • Implemented persisted checkpoint state for bounded source scans
  • Keep existence as the simplest invariant

Milestone 4: Failure Taxonomy Tests

  • Implemented source unavailable checker failure -> CheckFailed
  • Implemented verification-still-drifting repair failure -> RepairFailed
  • Next: target unavailable
  • Next: checker timeout
  • Next: repair Job crash before report publication
  • Next: repair retry/idempotency
  • Next: stale/corrupted target document

Milestone 5: Second Source/Derived Pair

  • Implemented Postgres -> Redis cache freshness
  • Added Redis cache indexer, Redis derived view CRD example, Redis Secret env resolution, Redis frontier evidence, and cache invalidation repair requests
  • Implemented Postgres -> ClickHouse analytics aggregate checks
  • Added ClickHouse analytics table initialization, backfill, drift demo, derived view CRD example, Secret env resolution, and backfill request repair mode

Milestone 6: Operator Shape

  • Define CRDs
  • Implemented first controller-runtime skeleton
  • Implemented synthetic Invariant reconciliation
  • Implemented job-backed Invariant reconciliation
  • Implemented checker Job creation and compact ConfigMap status handoff
  • Implemented scheduled checker Jobs through spec.checkIntervalSeconds
  • Implemented repair Job creation from explicitly allowed RepairPolicy
  • Implemented repair result ConfigMap handoff and verified status update
  • Added operator Dockerfile and in-cluster Deployment/RBAC sketch
  • Added Kubernetes-native demo stack for Postgres, Redpanda, OpenSearch, Redis, and ClickHouse
  • Converted demo data systems to StatefulSets with PVCs, probes, and resource requests/limits
  • Added Secret rotation fan-out to dependent Invariants
  • Enabled CRD status subresource
  • Update .status on Invariant
  • Run on kind/minikube

Milestone 7: Observability

  • Implemented Prometheus textfile metrics:
    • kubedataguard_drift_count
    • kubedataguard_checked_records
    • kubedataguard_counterexample_count
    • kubedataguard_slo_breach_count
    • kubedataguard_check_healthy
    • kubedataguard_cdc_frontier_status
  • Implemented OTEL-shaped JSONL span artifacts for check reports
  • OpenTelemetry traces around check and repair cycles
  • Grafana dashboard

Milestone 8: Advanced Checks

  • Implemented prefix-bucket checksum/Merkle-style narrowing
  • Schema compatibility checks
  • Sampled checks for large tables
  • Shard-aware checks

Stretch Ideas

  • Admission webhook that blocks deploying a pipeline without invariants
  • GitOps workflow for consistency SLOs
  • Slack/GitHub/Jira incident integration
  • Human approval for high-risk repairs
  • Multi-tenant invariant isolation
  • "Explain drift" page that shows likely root cause

About

Kubernetes-native data consistency SLO operator for derived data correctness

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages