A Python 3.11+ FastAPI service and Kafka worker for rule-based incident
analysis. The service exports traces and metrics with OpenTelemetry and adds
the active trace_id and span_id to application logs.
app/main.py- FastAPI entry point and telemetry lifecycleapp/api/analysis.py- REST analysis routesapp/domain/- domain modelsapp/service/- analysis logicapp/messaging/- Kafka contracts, producer, consumer, and workerapp/telemetry.py- OpenTelemetry configuration and worker instrumentstests/- unit and opt-in integration tests
pyproject.toml is the source of truth for Python dependencies.
python -m venv .venv
source .venv/bin/activate # Windows: .venv\Scripts\activate
python -m pip install .Start the API:
uvicorn app.main:app --reloadThe API is available at http://localhost:8000; its health endpoint is
GET http://localhost:8000/health.
The worker reuses the REST API's analyzer and runs as a separate process:
python -m app.messaging.workerKafka configuration is read from environment variables. The defaults target a local broker:
| Variable | Default |
|---|---|
KAFKA_BOOTSTRAP_SERVERS |
localhost:9092 |
ANALYSIS_REQUESTED_TOPIC |
incident.analysis.requested.v1 |
ANALYSIS_COMPLETED_TOPIC |
incident.analysis.completed.v1 |
ANALYSIS_FAILED_TOPIC |
incident.analysis.failed.v1 |
ANALYSIS_REQUEST_DLT_TOPIC |
incident.analysis.requested.v1.DLT |
ANALYSIS_WORKER_GROUP_ID |
incident-analyzer-workers-v1 |
AGENT_REQUESTED_TOPIC |
agent.execution.requested.v1 |
AGENT_EVENTS_TOPIC |
agent.execution.events.v1 |
AGENT_REQUEST_DLT_TOPIC |
agent.execution.requested.v1.DLT |
AGENT_WORKER_GROUP_ID |
incident-analyzer-agent-workers-v1 |
ANALYSIS_WORKER_MAX_ATTEMPTS |
3 |
ANALYSIS_WORKER_RETRY_BACKOFF_MS |
500 |
The worker consumes both analysis and V6 agent execution requests. It uses
manual offset commits and commits only after the completed, failed, or
dead-letter output has been acknowledged by Kafka. Agent step events are also
acknowledged before execution continues. Event IDs are deterministic across
Kafka redelivery so the Java audit consumer can process them idempotently.
Step, completed, and failed events include an eventType discriminator and
share the execution ID as their key, preserving per-execution Kafka ordering.
Agent execution requests carry the configured capability allowlist; the worker
will only plan and execute capabilities included in that request.
Each capability cycle emits PLAN, CAPABILITY_CALL, and OBSERVATION audit
steps before the terminal FINAL_RESULT step.
Both the API and worker initialize telemetry on startup and flush it during graceful shutdown. The API is instrumented with FastAPI instrumentation. The worker creates spans for Kafka consume, analysis, and publish operations and propagates W3C trace context through Kafka headers.
The V6 agent path adds agent.runtime.execute, agent.plan,
agent.capability.execute, agent.step.publish, and agent.result.publish
business spans around the existing Kafka spans.
The exporter uses OTLP/gRPC. Its default destination is
http://localhost:4317. Configure a collector with the standard OpenTelemetry
environment variables:
| Variable | Default | Purpose |
|---|---|---|
OTEL_EXPORTER_OTLP_ENDPOINT |
http://localhost:4317 |
Shared OTLP/gRPC collector endpoint |
OTEL_EXPORTER_OTLP_TRACES_ENDPOINT |
shared endpoint | Trace-specific endpoint override |
OTEL_EXPORTER_OTLP_METRICS_ENDPOINT |
shared endpoint | Metric-specific endpoint override |
OTEL_EXPORTER_OTLP_INSECURE |
derived from endpoint | Set to true for a plaintext collector; the default HTTP endpoint is insecure |
OTEL_EXPORTER_OTLP_HEADERS |
empty | Authentication or collector headers |
OTEL_EXPORTER_OTLP_TIMEOUT |
SDK default | Export timeout in seconds |
OTEL_METRIC_EXPORT_INTERVAL |
SDK default | Metric export interval in milliseconds |
OTEL_TRACES_SAMPLER |
parentbased_always_on |
SDK trace sampler |
OTEL_TRACES_SAMPLER_ARG |
empty | Optional sampler argument |
OTEL_SDK_DISABLED |
false |
Disable telemetry when set to true |
DEPLOYMENT_ENVIRONMENT |
local |
deployment.environment resource attribute |
The service resource attributes are service.name=incident-analyzer and
service.version=v5.
Example for a local plaintext collector:
export OTEL_EXPORTER_OTLP_ENDPOINT=http://localhost:4317
export OTEL_EXPORTER_OTLP_INSECURE=true
export DEPLOYMENT_ENVIRONMENT=local
uvicorn app.main:app --reloadNotable worker metrics include request, completion, failure, analysis duration,
Kafka publish failure, queue wait, invalid queue timestamp, and dead-letter
counters. The V6 agent runtime also exports
agent_runtime_executions_total, agent_runtime_duration_seconds,
agent_runtime_steps_total, agent_runtime_failures_total,
agent_capability_executions_total, and
agent_capability_duration_seconds. Queue wait requires an offset-aware
requestedAt; naive timestamps are rejected by the event contract.
The image installs the project from pyproject.toml and starts the API by
default:
docker build -t incident-analyzer .
docker run --rm -p 8000:8000 \
-e OTEL_EXPORTER_OTLP_ENDPOINT=http://host.docker.internal:4317 \
-e OTEL_EXPORTER_OTLP_INSECURE=true \
incident-analyzerOverride the default command to run the Kafka worker:
docker run --rm \
-e KAFKA_BOOTSTRAP_SERVERS=host.docker.internal:9092 \
-e OTEL_EXPORTER_OTLP_ENDPOINT=http://host.docker.internal:4317 \
-e OTEL_EXPORTER_OTLP_INSECURE=true \
incident-analyzer python -m app.messaging.workerhost.docker.internal is appropriate for Docker Desktop. Use the collector and
Kafka service names instead when the containers share a Compose network.
Run the unit test suite:
python -m pytestThe real Kafka integration test is opt-in and requires the configured topics:
KAFKA_INTEGRATION_TEST=1 python -m pytest \
tests/messaging/test_kafka_integration.py