This page is the detailed source for runtime choices, topology boundaries, and current architecture claims. Start with the architecture walkthrough for guided topology diagrams or the components walkthrough for the responsibility map, and use engineering status for the evidence and acceptance state behind those claims.
AgentFlow is an event-native metrics layer: business metrics are generated by operational events and stay live — the serving cache is invalidated when events arrive, with a measured 3.02 s p50 / 5.70 s p95 event-to-metric delay on the real Kafka→Flink→bridge path (2026-07-09 S8 snapshot) and 1.06 s p50 / 1.99 s p95 on the in-process demo shortcut, which skips Kafka/Flink (2026-06-06 demo snapshot). Each metric declares which events move it through versioned contracts, so the event→metric graph is a tested artifact rather than tribal knowledge.
Consumers are anything that needs the current number at the moment of decision: humans, dashboards, downstream services, and AI agents. Agents stress the same properties hardest — seconds-level freshness, semantic context, machine-readable contracts — but they are one consumer of the boundary, not its definition.
┌──────────────┐ ┌──────────────┐ ┌──────────────┐
│ Web Store │ │ Payment │ │ Inventory │
│ (events) │ │ Gateway │ │ Service │
└──────┬───────┘ └──────┬───────┘ └──────┬───────┘
│ │ │
▼ ▼ ▼
┌─────────────────────────────────────────────────────────┐
│ AgentFlow Platform │
│ │
│ Kafka → PyFlink → validated → Iceberg + ClickHouse → API│
└─────────────────────────────┬───────────────────────────┘
│
┌─────────▼──────────┐
│ Consumers │
│ (humans, dashboards│
│ services, agents) │
└────────────────────┘
- Streaming-first: Batch is a special case of streaming (bounded stream). One codebase, one semantics.
- Quality gates before storage: Bad data never reaches the serving layer. Agents never see it.
- Semantic over raw: Agents query entities and metrics, not tables and columns.
- Cost-aware: Every component has autoscaling and lifecycle policies. We measure $/GB processed.
- Observable: If you can't measure it, you can't operate it. Every stage emits latency, throughput, and error metrics.
The golden production topology ADR
selects one artifact and runtime: containerized PyFlink 2.3 on Kubernetes. Its
version and acceptance state are recorded in the
machine-readable project claims. The verified
streaming path today is Kafka → PyFlink 2.3 → events.validated → bridge →
ClickHouse → API, measured at 3.02 s p50 / 5.70 s p95. The accepted topology
adds an independently replayable Iceberg materializer from events.validated;
it is not a production claim until the clean-checkout acceptance gate passes.
- Ingestion: Events arrive via Kafka producers (orders, payments, clicks) or Debezium CDC connectors running on Kafka Connect.
- Processing: containerized PyFlink 2.3 validates, enriches, deduplicates, and routes events to
events.validatedorevents.deadletter. - Lake materialization: a dedicated consumer writes validated events to Iceberg. Production acceptance requires this consumer and its replay proof; the current verified benchmark did not include this hop.
- Quality: pre-storage gates check schema and semantic rules. Failures go to the dead-letter topic.
- Serving materialization: the bridge applies
events.validatedto ClickHouse idempotently by(tenant_id, event_id). After a successful apply it also pushes metric-cache invalidation (Redis channel + in-process callback); a journal-scan fallback covers writers that do not push. See Serving Bridge. - Serving: Agent API reads ClickHouse in the demo/production profile (
config/serving.yaml), with DuckDB as the local-dev/test compatibility store.
For CDC sources, Debezium/Kafka Connect handles source capture while a shared normalizer converts Postgres/MySQL envelopes into one canonical AgentFlow CDC contract before validation. See ADR 0005.
Same pipeline logic, no infrastructure dependencies:
- Generate:
local_pipeline.pycreates realistic e-commerce events - Validate: Schema validation (Pydantic) + semantic validation (business rules)
- Enrich: Domain enrichment per event type (order sizing, click classification, payment risk)
- Store: Validated events written to DuckDB, mirrored to the ClickHouse serving tables when that profile is up (
src/agentflow_runtime/processing/clickhouse_sink.py), and written to Iceberg via PyIceberg - Serve: Agent API reads from the configured serving backend (ClickHouse on the demo profile, DuckDB under
SERVING_BACKEND=duckdb) while/v1/healthreports Iceberg row counts - Catalog: Development uses a MinIO-backed REST catalog from
docker-compose.iceberg.yml(writing to the sameagentflow-lakeS3 object store as the Flink stack); production uses AWS Glue
Both paths use the same validator and enrichment code (src/agentflow_runtime/quality/, src/agentflow_runtime/processing/transformations/).
- Orchestration: Dagster triggers compaction, aggregation, and quality reports
- User profiles: Materialized from
orders_v2→users_enriched - Quality report: Row counts, null rates, dead letter ratio
- Compaction: Iceberg snapshot expiry + data file compaction (production only)
The serving layer has grown beyond the original read-only surface. The current API groups into four slices:
- Core agent reads:
/v1/entity,/v1/metrics,/v1/query,/v1/catalog,/v1/health - Discovery and audit:
/v1/search,/v1/contracts,/v1/lineage,/v1/changelog - Operational workflows:
/v1/batch,/v1/stream/events,/v1/deadletter,/v1/webhooks,/v1/alerts,/v1/slo - SDK contract: Python sync/async clients and a TypeScript client wrap the same HTTP surface
The API process also starts several background components:
- DuckDBPool: shared read cursors plus serialized writes for the local serving path
- QueryCache: Redis-backed metric cache with invalidation on new events
- WebhookDispatcher: polls validated pipeline events and delivers signed webhook callbacks
- AlertDispatcher: evaluates metric thresholds and records alert history
- OutboxProcessor: retries replay delivery from DuckDB outbox rows to Kafka
These components keep the local demo and the production architecture aligned around one agent-facing contract, even when the backing infrastructure differs.
See Architecture Decision Records for detailed trade-off analysis.
| Component | Choice | Runner-up | Key differentiator |
|---|---|---|---|
| Streaming | Kafka 3.7 (KRaft) | Pulsar | Ecosystem maturity, MSK managed service |
| CDC capture | Debezium + Kafka Connect | Python-native connectors | Mature Postgres/MySQL CDC, built-in offsets/schema history, one ops model |
| Processing | Flink 2.3 | Spark Structured Streaming | True event-time, lower latency, native watermarks |
| Storage | Iceberg 1.5 | Delta Lake | Vendor-neutral, hidden partitioning, time-travel |
| Local query | DuckDB | SQLite | Columnar, fast analytics, Iceberg support |
| Orchestration | Dagster | Airflow | Software-defined assets, better testing, type safety |
| API | FastAPI | Flask | Async, auto-docs, Pydantic integration |
| IaC | Terraform | Pulumi | Team familiarity, HCL readability, module ecosystem |
| Failure | Impact | Mitigation |
|---|---|---|
| Kafka broker down | Reduced throughput | 3-broker cluster, replication factor 3, min.insync.replicas=2 |
| Flink job crash | Processing stops | Exactly-once checkpointing (30s), auto-restart on failure |
| Bad data in source | Incorrect agent answers | Pre-storage quality gates, dead letter topic, alerting |
| S3 outage | No new data in serving | Flink checkpoints to S3 pause; resumes on recovery |
| API overload | Agent queries fail | Per-key rate limiting, DuckDB connection pooling, horizontal scaling |
| Redis unavailable | Cache or rate-limit degradation | Metric cache misses fall back to source queries; the rate limiter fails closed to a per-process cap so an outage cannot disable limiting fleet-wide |
| Webhook target failure | Lost downstream notification | Signed delivery logs, retries with backoff, alert/webhook history in DuckDB |
- API authentication: API key via
X-API-Keyheader (setAGENTFLOW_API_KEYSenv var). A key configured without akey_id(every environment key, and a key-file entry without one) gets an id derived from its key-lookup digest, stable across reloads, restarts and replicas (except a legacy hash-only entry, with neither a plaintext key nor akey_lookup: it has nothing to derive from and keeps a random id, which in a key file the process cannot write changes on every load); changing the pepper changes only an id derived from a plaintext key and never written back (every environment key, and a plaintext key-file entry in a file the process cannot write), while a writable key file keeps the id written to it on the first load and an entry with a storedkey_lookupkeeps the id derived from that digest. Requirements for a key's id: API key identity - Transport gate (audit P2-3):
AGENTFLOW_PROFILE=productionrefuses to boot over plaintext transport to an external ClickHouse/Redis/PostgreSQL (loopback exempt; deliberate exceptions named inAGENTFLOW_INSECURE_TRANSPORT_OK), and refuses a wildcard CORS origin outside demo mode. The ClickHouse client supports HTTPS with hostname verification and a private-CA bundle (CLICKHOUSE_SECURE,CLICKHOUSE_CA_CERT) - Rate limiting: Per-key sliding window with Redis backing when available and in-memory fallback for local/test, configurable via
AGENTFLOW_RATE_LIMIT_RPM(default: 120/min). Requirements: API key rate limiting - Health/docs exempt:
/v1/health,/docs,/metricsdon't require auth - No secrets in code: All credentials via environment variables
- Terraform state: Encrypted S3 backend with DynamoDB locking
- Kafka: TLS in-transit, SASL authentication (MSK config)
- S3: SSE-KMS encryption, bucket policy restricts to VPC endpoints
- Network: Private subnets, security groups per component
- Metrics: Prometheus scrapes
/metrics; Grafana dashboards cover pipeline health plus support, ops, and merch journeys. - Tracing: OpenTelemetry spans export to Jaeger through
OTEL_EXPORTER_OTLP_ENDPOINT; the production-like compose stack exposes Jaeger on:16686. - Logs: Structlog emits JSON logs with
trace_id,span_id,correlation_id, and tenant context so incidents can be traced across API, cache, and background loops. - Operational APIs:
/v1/alerts,/v1/webhooks,/v1/deadletter,/v1/slo, and/v1/stream/eventsare part of the control plane, not side tooling.
| Environment | Primary components | Purpose |
|---|---|---|
| Local demo | make demo: Redis + ClickHouse (Docker) + agentflow_runtime.processing.local_pipeline + FastAPI |
Fastest path for developers and SDK examples |
| Production-shaped local Docker | docker-compose.prod.yml with Kafka, Redis, Jaeger, Prometheus, Alertmanager, Grafana, API, and optional ClickHouse |
Observability and production-shaped debugging against a realistic local stack. A demo: no TLS, dev credentials, demo-mode auth |
| Lite E2E Docker | docker-compose.e2e.yml with the narrowed CI service set |
Faster E2E and smoke coverage without the full observability stack |
| Chaos harness | docker-compose.chaos.yml + Toxiproxy + pytest chaos suite |
Validate graceful degradation under Kafka/Redis failures |
| kind staging | helm/agentflow, k8s/, scripts/k8s_staging_up.sh |
Production-shaped staging on a local Kubernetes cluster |
| Production candidate | Kubernetes + containerized PyFlink 2.3 + Kafka + Iceberg materializer + ClickHouse + API | Accepted golden topology; not production-ready until the ADR acceptance gate passes |
Serving engine decision — fixed on ClickHouse (ADR 0006, ADR 0007). The serving engine is now a recorded decision, not an open question: the demo/production serving path fixes on ClickHouse, with DuckDB demoted to the local-dev / test and compatibility store (still first-class for
pytestand offline work — see the DuckDB pin intests/conftest.py).Phase 1 is executed (2026-07-02).
config/serving.yamldefaults tobackend: clickhouse;make demoanddocker-compose.prod.ymlbring up the ClickHouse service by default. The local pipeline mirrors serving-table writes to ClickHouse (src/agentflow_runtime/processing/clickhouse_sink.py), and the freshness-critical event scan (webhooks, metric-cache invalidation, SSE) goes through the serving backend (QueryEngine.fetch_pipeline_events) — so the event→metric axis works on the shipped engine, across process boundaries. Upserts are modeled as ReplacingMergeTree row versions withfinal=1reads. The Helm chart keeps the safe single-node DuckDB profile as default (it ships no ClickHouse service);serving.backend=clickhousewires an external service.Horizontal scaling stays gated (ADR 0009). The control plane (webhook queue, alert history, outbox, usage — plus webhook registrations and alert rules/runtime state in per-pod YAML files) is embedded per-pod state; scaling requires externalizing it, not only the serving engine. The externalization is now a recorded decision (ADR 0010): PostgreSQL behind a
ControlPlaneStoreport with the embedded store staying the default single-replica profile, rolled out in staged slices. The chart enforces the gate at render time — any multi-replica render fails untilcontrolPlane.store=postgresships and is set alongsideserving.backend=clickhouse.PII is not a serving-tier concern (2026-07-01). The demo serving warehouse holds no PII —
users_enriched/orders_v2carry only analytics columns — so the interim NL→SQL PII deny-gate and the entity redactor were guarding a surface that never exists in the demo, and both have been removed (see CHANGELOG). Real contact PII lives only in the DV2 business vault; its governance is engine-side there, not a dialect-pinned string parse in the serving tier. Executed 2026-07-02 (ADR 0006 Phase 2):warehouse/agentflow/dv2/governance/— fail-closed column grants fordv2_analyst, per-jurisdictiondv2_pii_officer__<branch>roles, row policies onhub_customer,SQL SECURITY DEFINERMDM views, PII-freecustomer_360; verified live, 29/29 probes on the current script (docs/perf/vault-pii-governance-verify-2026-07-03.md).sql_guardremains, scoped to what it can actually enforce: SELECT-only, no DML, the tenant table allow-list, and the recursive-CTE shadow reject. Tracked inroad-to-9.8.md.Operational surfaces have a recorded serving split (ADR 0011, 2026-07-03). The planned ops layer (Order 360 timeline, stuck-orders worklist, exception inbox — the workflows of
domain.md§4) composes exactly the two existing ports: analytical reads (entity point-reads, thepipeline_eventsjournal, open-order scans) on the serving backend, transactional triage state (dead-letter lifecycle, the exception-triage overlay, the manual-work counter) on theControlPlaneStore— no third data path. Endpoint contracts, the SLA stage model, and exception sources are pinned indocs/ops-surfaces-spec.md; implementation is staged (D2–D4).A three-node demo topology is designed (ADR 0012, 2026-07-04). The public Hugging Face demo grows from one container to a center (
mskHQ) plus two edge branches (spb/ekb—domain.md§1) that push operational events over HTTPS, so the artifact shows the branch-distribution story live. This is a distinct axis from horizontal scaling: the three nodes are separate single-replica Spaces with embedded control planes — not the ADR 0010 PostgreSQL scale profile — HTTPS substitutes for the Kafka→Flink transport, and the federated live layer is ephemeral on a deterministic re-seeded baseline. Node roles, the ingest contract, sleep choreography, and the N1–N12 test invariants are pinned indocs/architecture/three-node-demo-topology.md; implementation is F2 and the Space deploy is an owner gate.
| Capability | Implementation | Architectural impact |
|---|---|---|
| Durable replay and outbox | OutboxProcessor, dead-letter replay, DuckDB outbox rows |
Background delivery is retried without coupling API latency to Kafka availability |
| Redis-backed rate limiting | AuthManager + RateLimiter with a fail-closed per-process fallback cap |
Per-key throttling stays centralized in prod, while local/test continues (capped per process) when Redis is absent |
| Typed SDK surface | sdk/ and sdk-ts/ |
Python and TypeScript agents consume one HTTP contract instead of bespoke adapters |
| Distributed observability | OTel tracing, structlog correlation, /metrics, Grafana, Jaeger |
API, background jobs, and streaming workflows share the same debugging context |
| Chaos engineering | tests/chaos/, docker-compose.chaos.yml, config/toxiproxy.json |
Failure handling is exercised continuously rather than assumed from code review |
| Kubernetes staging | helm/agentflow, k8s/kind-config.yaml, staging scripts |
One signed/attested workflow digest is pulled through Helm and held to smoke/E2E validation before production |
| DevContainer DX | .devcontainer/ with Docker-in-Docker, Helm/kubectl, kind, toxiproxy-cli |
Contributors get one workspace that can run local demo, chaos, and staging workflows |