Skip to content

Latest commit

 

History

History
246 lines (206 loc) · 19.7 KB

File metadata and controls

246 lines (206 loc) · 19.7 KB

Architecture Overview

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.

Context

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.

System Context (C4 Level 1)

┌──────────────┐     ┌──────────────┐     ┌──────────────┐
│  Web Store   │     │   Payment    │     │  Inventory   │
│  (events)    │     │   Gateway    │     │   Service    │
└──────┬───────┘     └──────┬───────┘     └──────┬───────┘
       │                    │                    │
       ▼                    ▼                    ▼
┌─────────────────────────────────────────────────────────┐
│                    AgentFlow Platform                    │
│                                                         │
│ Kafka → PyFlink → validated → Iceberg + ClickHouse → API│
└─────────────────────────────┬───────────────────────────┘
                              │
                    ┌─────────▼──────────┐
                    │     Consumers      │
                    │ (humans, dashboards│
                    │  services, agents) │
                    └────────────────────┘

Key Design Principles

  1. Streaming-first: Batch is a special case of streaming (bounded stream). One codebase, one semantics.
  2. Quality gates before storage: Bad data never reaches the serving layer. Agents never see it.
  3. Semantic over raw: Agents query entities and metrics, not tables and columns.
  4. Cost-aware: Every component has autoscaling and lifecycle policies. We measure $/GB processed.
  5. Observable: If you can't measure it, you can't operate it. Every stage emits latency, throughput, and error metrics.

Data Flow

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.

Golden container path

  1. Ingestion: Events arrive via Kafka producers (orders, payments, clicks) or Debezium CDC connectors running on Kafka Connect.
  2. Processing: containerized PyFlink 2.3 validates, enriches, deduplicates, and routes events to events.validated or events.deadletter.
  3. 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.
  4. Quality: pre-storage gates check schema and semantic rules. Failures go to the dead-letter topic.
  5. Serving materialization: the bridge applies events.validated to 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.
  6. 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.

Local: Generate → Validate → Enrich → DuckDB + Iceberg

Same pipeline logic, no infrastructure dependencies:

  1. Generate: local_pipeline.py creates realistic e-commerce events
  2. Validate: Schema validation (Pydantic) + semantic validation (business rules)
  3. Enrich: Domain enrichment per event type (order sizing, click classification, payment risk)
  4. 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
  5. Serve: Agent API reads from the configured serving backend (ClickHouse on the demo profile, DuckDB under SERVING_BACKEND=duckdb) while /v1/health reports Iceberg row counts
  6. Catalog: Development uses a MinIO-backed REST catalog from docker-compose.iceberg.yml (writing to the same agentflow-lake S3 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/).

Batch Path (daily)

  1. Orchestration: Dagster triggers compaction, aggregation, and quality reports
  2. User profiles: Materialized from orders_v2 → users_enriched
  3. Quality report: Row counts, null rates, dead letter ratio
  4. Compaction: Iceberg snapshot expiry + data file compaction (production only)

Serving & Control Plane

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.

Technology Choices

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 Modes & Mitigations

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

Security

Implemented

  • API authentication: API key via X-API-Key header (set AGENTFLOW_API_KEYS env var). A key configured without a key_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 a key_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 stored key_lookup keeps the id derived from that digest. Requirements for a key's id: API key identity
  • Transport gate (audit P2-3): AGENTFLOW_PROFILE=production refuses to boot over plaintext transport to an external ClickHouse/Redis/PostgreSQL (loopback exempt; deliberate exceptions named in AGENTFLOW_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, /metrics don't require auth
  • No secrets in code: All credentials via environment variables
  • Terraform state: Encrypted S3 backend with DynamoDB locking

Production (via infrastructure)

  • 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

Observability & Operations

  • 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/events are part of the control plane, not side tooling.

Deployment Topologies

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 pytest and offline work — see the DuckDB pin in tests/conftest.py).

Phase 1 is executed (2026-07-02). config/serving.yaml defaults to backend: clickhouse; make demo and docker-compose.prod.yml bring 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 with final=1 reads. The Helm chart keeps the safe single-node DuckDB profile as default (it ships no ClickHouse service); serving.backend=clickhouse wires 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 ControlPlaneStore port 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 until controlPlane.store=postgres ships and is set alongside serving.backend=clickhouse.

PII is not a serving-tier concern (2026-07-01). The demo serving warehouse holds no PII — users_enriched/orders_v2 carry 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 for dv2_analyst, per-jurisdiction dv2_pii_officer__<branch> roles, row policies on hub_customer, SQL SECURITY DEFINER MDM views, PII-free customer_360; verified live, 29/29 probes on the current script (docs/perf/vault-pii-governance-verify-2026-07-03.md). sql_guard remains, scoped to what it can actually enforce: SELECT-only, no DML, the tenant table allow-list, and the recursive-CTE shadow reject. Tracked in road-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, the pipeline_events journal, open-order scans) on the serving backend, transactional triage state (dead-letter lifecycle, the exception-triage overlay, the manual-work counter) on the ControlPlaneStore — no third data path. Endpoint contracts, the SLA stage model, and exception sources are pinned in docs/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 (msk HQ) 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 in docs/architecture/three-node-demo-topology.md; implementation is F2 and the Space deploy is an owner gate.

v1-v6 Capability Map

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