Stateful Eval Worker
In-process stateful evaluation worker — runs temporal metric computation directly as a library, bypassing the HTTP service.
Source: workers/stateful-eval-worker
Overview
The Stateful Eval worker scores agent sessions with the Stateful Evals
library. It takes a session ID off a RabbitMQ queue, loads the session's spans from ClickHouse
through the OXP API library (oxp.client.local.LocalClient), and evaluates them in-process with
stateful_evals_be.TemporalMetricsProcessor. It writes per-span Groundedness, IntentRecognition,
and Relevancy scores and a session trajectory_score to the Neo4j knowledge graph, then forwards
the message. Evaluation makes model calls, which can cost money.
StatefulEvalWorker (src/stateful_eval_worker/worker.py) subclasses the shared
worker-base BaseWorker.
This repository does not include an upstream worker that routes sessions to this queue. What's certain from this worker's own source is its queue name:
- Input queue:
new_session_to_stateful_eval(queues.NEW_SESSION_TO_STATEFUL_EVAL_Q) - Output queue(s): none by default — configurable, comma-separated
- CLI:
stateful-eval-worker-cli
How it works
For each message, StatefulEvalWorker.handle_message:
- Fetches every span of
session_idwithstateful_evals_be.integrations.oxp.fetch_session_spans(500 per page, oldest first, converted byoxp_span_to_otel). - Builds a new
TemporalMetricsProcessorand callsevaluate_spans(spans, session_id=...)— the same processorstateful_evals_be.evaluate_spansuses. No policy override is passed, so the policy is whatever system message is found in the spans. - Writes the metrics to Neo4j when
PUSH_METRICSis on andresult.erroris empty. - Publishes one output message to each configured output queue.
Fetching, evaluation, and metric writes run in threads; up to MAX_INFLIGHT_MESSAGES messages run
at once, sharing one LocalClient. If the processor reports an error (e.g. "No spans found for
session"), the worker writes no metrics but still publishes. If fetching or writing raises, it logs
the error, sends an error feedback event, and publishes nothing. Either way the message is
acknowledged, not requeued.
What it writes
The worker only reads ClickHouse; it writes to Neo4j via LocalClient.write_span_metrics and
write_session_metrics:
| Metric | Attached to | Value |
|---|---|---|
mdt.Groundedness, mdt.IntentRecognition, mdt.Relevancy |
(:Session {sessionId})-[:hasSpan]->(:Span {spanId}), created if missing |
0 or 1 per judged llm or tool span; tool spans get no Groundedness |
trajectory_score |
An existing Session node with that sessionId; nothing is written if there is none |
0 or 1; reasoning is trajectory_reasoning |
Each metric is a Metric node linked by hasMetric, with provider stateful_evals, source
StatefulEval, metricId equal to the name, and reasoning. It's matched on metricName and
resourceId, so a re-run overwrites it. The OXP API's quality category reads trajectory_score
as a source of MDT (trajectoryQuality) in api/oxp/core/config.py.
Usage
# Liveness check: prints "I'm Alive"
uv run --package stateful-eval-worker stateful-eval-worker-cli --test
# Queue mode
uv run --package stateful-eval-worker stateful-eval-worker-cli \
--rabbitmq-url amqp://<user>:<password>@localhost:5672/ --output-queue <next-queue>
# Run-once mode (Argo Workflows compatible): reads --input, writes --output, no broker connection
echo '{"session_id": "<session-id>"}' > /tmp/stateful-eval-input.json
uv run --package stateful-eval-worker stateful-eval-worker-cli --run-once \
--rabbitmq-url amqp://<user>:<password>@localhost:5672/ \
--input /tmp/stateful-eval-input.json --output /tmp/stateful-eval-output.json
--run-once writes {"workflow_id": ..., "session_id": ...} to --output, generating a
workflow_id if the input has none. --input also accepts inline JSON; a JSON array is processed
as a batch and written back as an array. In queue mode, --max-sessions N caps the number of
messages processed (-1, the default, means no limit). Both modes read spans from ClickHouse —
the worker has no file input for sessions. To evaluate a sample trajectory without a database, use
stateful_evals_be.evaluate_file directly, as the
stateful-evals quickstart does.
Configuration
Flags override environment variables, except MAX_INFLIGHT_MESSAGES, which overrides
--max-inflight-messages. From a source checkout, cli.py also loads
workers/stateful-eval-worker/.env without overriding variables already set.
| Env var | CLI flag | Default | Purpose |
|---|---|---|---|
RABBITMQ_URL |
--rabbitmq-url |
none | Broker URL; this or the split variables below are required, even with --run-once |
RABBITMQ_HOST / RABBITMQ_USER / RABBITMQ_PASSWORD / RABBITMQ_PORT / RABBITMQ_VHOST |
(env only) | none / none / none / 5672 / / |
Alternative to RABBITMQ_URL; host, user, and password are required |
STATEFUL_EVAL_INPUT_QUEUE |
--input-queue |
new_session_to_stateful_eval |
Input queue |
STATEFUL_EVAL_OUTPUT_QUEUE |
--output-queue (repeatable) |
none | Output queue(s), comma-separated |
STATEFUL_EVAL_FEEDBACK_QUEUE |
--feedback-queue |
none | Feedback queue |
MAX_INFLIGHT_MESSAGES |
--max-inflight-messages |
32 |
Concurrent messages / prefetch count; -1 uses the CPU count |
OPENAI_API_KEY |
--llm-api-key |
empty | Judge model API key |
LLM_MODEL_NAME |
--llm-model-name |
gpt-4o |
Judge model |
LLM_BASE_MODEL_URL_MCE |
--llm-base-model-url |
https://api.openai.com/v1 |
Judge model endpoint |
SAMPLING_STRATEGY, SAMPLING_EARLY_RATE, SAMPLING_MID_RATE |
(env only) | none / 0.25 / 0.60 |
tail_weighted judges these shares of spans in the first 40% and next 30%; the rest and anchor spans are always judged |
PUSH_METRICS |
(env only) | true |
false/0/no skips all Neo4j writes |
CLICKHOUSE_HOST / CLICKHOUSE_PORT / CLICKHOUSE_USERNAME / CLICKHOUSE_PASSWORD / CLICKHOUSE_DATABASE |
(env only) | localhost / 8123 / admin / admin / default |
Span store (raw OTel traces) |
NEO4J_HOST / NEO4J_PORT / NEO4J_USERNAME / NEO4J_PASSWORD / NEO4J_DATABASE |
(env only) | localhost / 7687 / neo4j / empty / neo4j |
Knowledge graph, used when PUSH_METRICS is on; NEO4J_AUTH=user/password also sets credentials |
ClickHouse and Neo4j defaults come from oxp.core.config.settings in the api
package. LocalClient.from_settings connects to bolt://<NEO4J_HOST>:<NEO4J_PORT> unless
NEO4J_HOST is already a URI — a bare host must not include a port. The worker-base Neo4j
connector created at startup isn't used by this worker; if it fails to connect, it only logs a
warning.
Build & Run
# Build a wheel (run from the monorepo root)
uv build workers/stateful-eval-worker
# Run the CLI directly
uv run --package stateful-eval-worker stateful-eval-worker-cli --test
Docker
docker build -f workers/stateful-eval-worker/deploy/docker/Dockerfile \
-t stateful-eval-worker .
docker run --rm stateful-eval-worker --test
entrypoint.sh forwards the container arguments to stateful-eval-worker-cli (or the
INPUT_WORKFLOW string when none are given) and logs the --input file. docker compose up
stateful-eval also starts RabbitMQ and Neo4j; ClickHouse isn't a compose service, so
CLICKHOUSE_HOST defaults to host.docker.internal.
Development
uv run --package stateful-eval-worker --extra dev pytest workers/stateful-eval-worker/tests -v
The tests mock span fetching, the processor, and the OXP client — no broker, database, or model is needed.