This page owns the exact local commands, supported runtime settings, and recurring maintenance for the compose stack, local Flink, and DuckDB paths. It is for the local and demo operator who is bringing those services up, inspecting them, or recovering a local incident. Production-incident response lives in the on-call runbooks; controlled gates and specialized procedures live in the operations index.
Audience: local/demo operator for the compose stack, local Flink, and DuckDB paths
Prerequisites: Docker Compose, DuckDB CLI, Kafka CLI tools; kind staging also requires docker, kubectl, helm, and kind as named by scripts/k8s_staging_up.sh
| Service | Local URL | Health check |
|---|---|---|
| Agent API | http://localhost:8000/docs | curl http://localhost:8000/v1/health |
| Kafka | localhost:9092 | kafka-broker-api-versions --bootstrap-server localhost:9092 |
| Flink UI | http://localhost:8081 | curl http://localhost:8081/overview |
| MinIO | http://localhost:9001 | curl http://localhost:9000/minio/health/live |
| Prometheus | http://localhost:9090 | curl http://localhost:9090/-/healthy |
| Grafana | http://localhost:3000 | curl http://localhost:3000/api/health |
| Jaeger | http://localhost:16686 | curl -I http://localhost:16686 |
| Toxiproxy API | http://localhost:8474 | curl http://localhost:8474/proxies |
The troubleshooting walkthrough helps classify a symptom and select its owning procedure. The observability walkthrough explains how to combine these signals. This runbook owns their exact local commands, supported runtime settings, and recurring maintenance.
Scrape the API metrics directly:
curl http://localhost:8000/metricsThe production-shaped compose stack wires those metrics to Prometheus and Grafana at the URLs in the quick-reference table. Use Jaeger from the same table to inspect traces for a failing or slow request.
The supported OpenTelemetry settings are:
OTEL_EXPORTER_OTLP_ENDPOINT=http://jaeger:4317
OTEL_SERVICE_NAME=agentflow-api
OTEL_SDK_DISABLED=true # supported way to run without tracingOpenTelemetry is a mandatory runtime dependency and the API imports its
telemetry setup outright. Disabling tracing is an explicit configuration
choice through OTEL_SDK_DISABLED=true; an unimportable telemetry module
fails boot instead of silently degrading to no tracing (audit F-13).
Query analytics keeps a peppered fingerprint of each query question, not the
question, unless an operator opts in to a redacted copy. Retention defaults to
30 days. The Helm chart provides a disabled-by-default CronJob that calls the
authenticated API; it never opens DuckDB directly or mounts the API PVC. Enable
it with an operator-managed Secret containing admin-key:
analyticsRetention:
enabled: true
schedule: "0 3 * * *"
retentionDays: "" # use AGENTFLOW_QUERY_ANALYTICS_RETENTION_DAYS, default 30
dryRun: false
concurrencyPolicy: ForbidThe production values contract requires the job to be enabled with
dryRun: false and concurrencyPolicy: Forbid. Validate a new schedule in a
non-production values file with dryRun: true, then switch to the final mode.
For a manual preview, call the same API; a dry run does not call the store's
prune operation:
curl -X POST \
-H "X-Admin-Key: <admin-key>" \
-H "Content-Type: application/json" \
-d '{"retention_days":30,"dry_run":true}' \
http://localhost:8000/v1/admin/analytics/retentionscripts/prune_query_analytics.py is retained for offline maintenance only;
do not schedule it beside a live API process against the same DuckDB/PVC. See the
security policy
before retention or tenant-erasure work.
For an optional service-backed path, verify the engine and inspect the default stack before changing data or ports:
docker version
docker compose version
docker compose ps
docker compose logs kafka flink-jobmanagerIf the expected local store is missing or surprising, confirm the selected path and visible files:
=== "macOS / Linux"
```bash
echo "$DUCKDB_PATH"
ls *.duckdb*
```
=== "PowerShell"
```powershell
Write-Output $env:DUCKDB_PATH
Get-ChildItem -Filter "*.duckdb*"
```
When the default API address is already owned, keep the existing process and start the walkthrough API on another port:
uvicorn agentflow_runtime.serving.api.main:app --host 0.0.0.0 --port 8001Only after stale named volumes are the confirmed cause and their data is disposable, reset the default stack. This deletes those local volumes:
docker compose down -vmake demo # Starts Redis + ClickHouse, provisions the serving store,
# seeds 500 events through the pipeline, starts the APIThe API does not create or seed a serving store on boot — a serving identity
should not hold CREATE/ALTER/INSERT, and an empty production store should not be
handed demo rows for being empty (audit P0-2). make demo therefore runs the
provisioning step explicitly:
python -m agentflow_runtime.serving.provision --schema --seed # demo bring-up
python -m agentflow_runtime.serving.provision --schema # what a real deployment runsBoth are idempotent and re-runnable, and they provision every store the API
reads: on the ClickHouse profile that is ClickHouse and the embedded DuckDB,
which still carries the control-plane state (webhook queue, dead-letter inbox).
Compose runs it as the serving-init one-shot service; Helm runs it as a
pre-install/pre-upgrade Job.
make pipeline # 10 events/sec, writes to agentflow_demo.duckdbduckdb agentflow_demo.duckdbSELECT COUNT(*) FROM pipeline_events;
SELECT * FROM orders_v2 ORDER BY created_at DESC LIMIT 5;
SELECT topic, COUNT(*) FROM pipeline_events GROUP BY topic;make clean # Removes .duckdb files + caches
make demo # Start freshdocker compose -f docker-compose.prod.yml up -d
python scripts/wait_for_services.py --url http://127.0.0.1:8000 --timeout 120Use this path for E2E checks and incidents that depend on Redis, Jaeger, Prometheus, or Grafana.
make flink-localThis builds the local Python 3.11 Flink image, starts the required Kafka, MinIO, and Flink services, and submits src/agentflow_runtime/processing/flink_jobs/stream_processor.py to the local cluster.
For session aggregation behavior, checkpoint tuning, Kubernetes Operator compatibility, and the isolated PyFlink environment, use the Flink operator reference.
Verify the run here:
- Flink Web UI: http://localhost:8081
- Valid events:
events.validated - Invalid events:
events.deadletter
docker compose -f docker-compose.yml -f docker-compose.cdc.yml build kafka-connect
docker compose -f docker-compose.yml -f docker-compose.cdc.yml up -d kafka cdc-kafka-init postgres-source mysql-source kafka-connect
docker compose -f docker-compose.yml -f docker-compose.cdc.yml run --rm cdc-register-connectorsVerify connector state:
curl -fsS http://localhost:8083/connectors/agentflow-postgres-cdc/status
curl -fsS http://localhost:8083/connectors/agentflow-mysql-cdc/statusRun the optional Docker CDC integration test against the running stack:
AGENTFLOW_RUN_CDC_DOCKER=1 python -m pytest -p no:schemathesis tests/integration/test_cdc_capture.py -qStop the local CDC stack when finished:
docker compose -f docker-compose.yml -f docker-compose.cdc.yml downKafka auto-create is disabled locally, so cdc-kafka-init pre-creates raw table topics, Debezium heartbeat topics, the MySQL signal topic cdc.mysql, and Kafka Connect internal topics. The MySQL schema history topic must use cleanup.policy=delete with unlimited retention; Debezium 3.5 JSON schema-history records can be keyless, and a compacted topic rejects those records.
For Kubernetes installs, choose one Kafka Connect source-credential mode:
- Demo/staging chart-managed credentials: keep
secrets.create=trueand leavesecrets.existingSecretempty. - Externally managed credentials: set
secrets.create=falseand setsecrets.existingSecretto a Kubernetes Secret that containspostgres.propertiesandmysql.properties.
The chart schema rejects mixed or missing modes so a rendered deployment cannot reference a missing source-credential Secret.
Production source attachment has a separate gate. Before creating any connector for a real Postgres or MySQL database, complete Production CDC Source Onboarding.
# List running jobs
curl http://localhost:8081/jobs
# Cancel a job (it will restart from last checkpoint)
curl -X PATCH http://localhost:8081/jobs/<job-id>?mode=cancel
# Resubmit
flink run -py src/agentflow_runtime/processing/flink_jobs/stream_processor.pykafka-consumer-groups --bootstrap-server localhost:9092 \
--group agentflow-stream-processor --describeIf lag > 100k: Flink is behind. Check Flink UI for backpressure indicators.
# Stop the Flink job first, then:
kafka-consumer-groups --bootstrap-server localhost:9092 \
--group agentflow-stream-processor \
--reset-offsets --to-latest --execute \
--topic orders.rawkafka-console-consumer --bootstrap-server localhost:9092 \
--topic events.deadletter --from-beginning --max-messages 10 | jq .The supported promotion path is the manual Staging Deploy workflow. Supply
the successful Container Attestation build run ID, its exact 40-character
main-branch source SHA, and confirmation PROMOTE. The workflow validates the
run, packet, cosign signature, and GitHub build provenance before it creates a
kind cluster.
For an isolated rehearsal on a Docker/kind host, first complete those same verification checks and download the selected packet. Then pass its verified digest-only values file explicitly:
PROMOTION_VALUES_FILE=/absolute/path/to/image-values.yaml bash scripts/k8s_staging_up.sh
bash scripts/k8s_staging_down.shscripts/k8s_staging_up.sh expects docker, kubectl, helm, and kind on
the path. It refuses a missing or non-regular promotion values file, pulls the
immutable API subject through Helm, and runs the smoke test. It does not build
or load an API image.
- Check
curl http://localhost:8000/v1/health. - If the endpoint times out, inspect
docker compose -f docker-compose.prod.yml ps. - Read the latest API logs:
docker compose -f docker-compose.prod.yml logs agentflow-api --tail 50. - Open Jaeger at
http://localhost:16686and look for hanginghttp.request,nl_to_sql, orduckdb.queryspans. - Restart only the API first:
docker compose -f docker-compose.prod.yml restart agentflow-api. - If the API still cannot become healthy, switch to the recovery steps in
docs/operations/disaster-recovery.md.
- Check Grafana and compare
/v1/healthfreshness against the expected window. - Measure Kafka lag with
kafka-consumer-groups --bootstrap-server localhost:9092 --group agentflow-stream-processor --describe. - If Kafka lag is high, inspect broker disk and partition skew before changing partition count.
- If Flink shows backpressure, inspect TaskManager memory and parallelism in the Flink UI.
- If the sink is slow, check Iceberg or object-store write latency and run the relevant compaction or snapshot cleanup job.
- Check Flink UI -> Jobs -> Failed -> Exception.
- Look for one of the common causes: OOM (
taskmanager.memory.process.size), checkpoint timeout (execution.checkpointing.timeout), or Kafka connectivity problems. - Let Flink restart from the last checkpoint once.
- If it fails again with the same exception, fix the root cause before forcing another restart.
- Check aggregate status:
curl -H "X-API-Key: <key>" http://localhost:8000/v1/deadletter/stats. - List the newest failures:
curl -H "X-API-Key: <key>" "http://localhost:8000/v1/deadletter?page=1&page_size=20". - Inspect one event in detail to confirm whether the issue is schema drift, semantic validation, or downstream delivery.
- If the payload is correctable, replay it through
POST /v1/deadletter/{event_id}/replay. - If the source is producing invalid payloads, stop replay attempts, notify the source owner, and add filtering or contract fixes first.
- List active rules with
curl -H "X-API-Key: <key>" http://localhost:8000/v1/alerts. - Inspect recent history with
curl -H "X-API-Key: <key>" http://localhost:8000/v1/alerts/<alert_id>/history. - Check whether the rule is flapping or escalating too aggressively in
config/alerts.yaml. - If the rule is noisy and unactionable, temporarily disable it with
DELETE /v1/alerts/{alert_id}or update the threshold and cooldown. - After the metric stabilizes, recreate or update the alert and confirm the webhook history has normalized.
- List registrations with
curl -H "X-API-Key: <key>" http://localhost:8000/v1/webhooks. - Inspect the failing webhook logs with
curl -H "X-API-Key: <key>" http://localhost:8000/v1/webhooks/<webhook_id>/logs. - Trigger a synthetic delivery through
POST /v1/webhooks/{webhook_id}/test. - If the failure happens only in kind staging, verify the host loopback relay configured by
scripts/k8s_staging_up.sh. - Once the target endpoint is healthy again, replay or re-trigger the expected event path.
- Check rotation state with
curl -H "X-Admin-Key: <admin-key>" http://localhost:8000/v1/admin/keys/<key_id>/rotation-status. - Verify
requests_on_old_key_last_hour; if it is non-zero, some client still uses the old secret. - Roll out the new key to every caller and wait until old-key traffic drops to zero.
- Revoke the old secret with
POST /v1/admin/keys/{key_id}/revoke-old. - Confirm the old key now returns
401and the new key still passes/v1/health.
- Check Kafka first:
kafka-broker-api-versions --bootstrap-server localhost:9092. - Confirm producers are still sending data and inspect their logs.
- Verify the Flink job is running and consuming from the expected topics.
- Check network reachability and DNS between Flink, Kafka, and the sink services.
- If you recently changed chaos settings, confirm Toxiproxy has been reset and all proxies are enabled.
For restore drills, backup verification, or host loss scenarios, use docs/operations/disaster-recovery.md.
- Review dead letter topic volume (should be < 0.1% of total)
- Check Kafka disk usage (alert at 80%)
- Review Grafana dashboards for anomalies
- Review noisy alerts and webhook delivery failures before they become chronic incidents
- Iceberg snapshot expiry:
CALL system.expire_snapshots('table', older_than => 30d) - Iceberg compaction:
CALL system.rewrite_data_files('table') - Review and rotate API keys
- Run
pytest tests/chaos/ -v --tb=shortagainst the current compose stack - Cost review: compare actual vs projected spend
The two Iceberg lines above are the only retention for table data. S3
lifecycle in infrastructure/terraform/modules/storage/main.tf deliberately
expires nothing under the warehouse prefix: object age says nothing about
which manifests still reference a file, so an S3 rule there deletes files out
from under live snapshots instead of cleaning the table. If table storage is
growing, the answer is expire_snapshots and rewrite_data_files, never a
new lifecycle rule.