Anomaly Detection Worker
Standalone anomaly detection worker: within one semantic group of sessions, it separates typical
sessions from outliers on three layers (text, execution graph, metrics) and writes an
AnomalyReport per layer.
Source: workers/anomaly-detection-worker
!!! info "Relationship to the analysis worker"
The default docker-compose deployment runs the combined
analysis worker, which embeds this same logic alongside consistency and
normal behaviour. This standalone worker exists for deployments that want to scale or schedule
anomaly detection separately (for example as an Argo Workflows step).
The key idea: anomalies here are relative to a peer group, never absolute. A session is an outlier because it differs from the other sessions doing the same kind of work, so the grouping stage must run first and the group must contain enough sessions for "typical" to mean anything.
Input
AnomalyDetectionInputMessage:
{
"session_id": "session-abc",
"group_id": "group-xyz",
"group_hash": "abc123",
"sessions": [
{
"session_id": "session-abc",
"metrics": {"Cost": 1.0, "Latency": 230},
"output_content": "...",
"output_embedding": [0.1, 0.2],
"execution_graph": {"nodes": [], "edges": []}
}
]
}
| Field | Required | Notes |
|---|---|---|
session_id |
yes | Used for the idempotency pre-check and the group-membership guard |
group_id |
yes | The SemanticGroup to analyse |
group_hash |
no | Defaults to "", which disables the hash guard |
sessions |
no | Defaults to [] |
sessions selects the data path:
- Non-empty — analysed as-is; Neo4j is never read, and no embedding model is needed.
- Empty — needs
--embedding-modeland a live DB connection. The worker waits for the hierarchical-grouping lock, runsanalysis_pre_check(group_id, "AnomalyReport", group_hash, session_id), loads the group withget_analysis_data_for_semantic_group, and drops the message ifsession_idis not actually a member of the group.
Output
Queue output: none. The CLI hardcodes output_queue=[], so this worker is terminal.
Graph output: one AnomalyReport per layer, linked from the SemanticGroup by
hasAnomalyReport:
| Property | Value |
|---|---|
id |
sha256(group_id + layer + metric_name) |
dataType |
text, graph, or metric |
inlierSessions / outlierSessions |
Session id lists |
reason |
Human-readable layer description |
scores |
JSON list of per-session model scores |
threshold |
Decision threshold |
nodeHash |
group_hash, or sha256 of the sorted session ids when group_hash is empty |
aboutMetric |
Set on metric reports only, linking to the Metric node |
The reason strings are fixed: Anomalies detected in text embeddings,
Anomalies detected in graph structure, and Anomalies detected in metric <name>.
Ingestion goes through verify_kg_object, so a report that violates the SHACL shapes raises
OntologyValidationError and is not written.
Note:
scoresandthresholdare always empty The detectors indembuild theirAnomalyDetectionResultfrom the model'sfit_predictlabels only; they never populatescoresorthreshold. Soscoresis persisted as[]andthreshold, beingNone, is written as0.0. Treat the inlier/outlier split as the real output and ignore these two fields.
File output (run-once): AnomalyDetectionOutputMessage, echoing the session data it analysed
plus an anomalies array:
{
"session_id": "session-abc",
"group_id": "group-xyz",
"sessions": [{"session_id": "session-abc", "...": "..."}],
"anomalies": [
{
"inlier_sessions": ["session-abc", "session-def"],
"outlier_sessions": ["session-ghi"],
"reason": "Anomalies detected in text embeddings",
"layer": "text",
"scores": [],
"metadata": {"model_name": "isolation_forest"}
}
]
}
Position in the pipeline
grouping-worker / hierarchical-grouping-worker
│
└──► new_session_to_anomaly ──► anomaly-detection-worker ──► Neo4j (AnomalyReport)
| Input queue | new_session_to_anomaly (ANOMALY_INPUT_QUEUE) |
| Output queue | none |
| CLI | anomaly-detection-worker-cli |
Warning: No producer publishes to this queue in the default deployment In
docker-compose.ymlthe grouping workers publish tonew_session_to_analysis, which the combined analysis worker consumes. Running this worker means either repointing a producer atnew_session_to_anomalyor driving it in run-once mode.
Quick start
# Liveness check
uv run --directory workers/anomaly-detection-worker anomaly-detection-worker-cli --test
Queue mode:
anomaly-detection-worker-cli \
--rabbitmq-url amqp://<user>:<password>@rabbitmq:5672/ \
--input-queue new_session_to_anomaly \
--embedding-model azure/text-embedding-3-small
Run-once against a file:
anomaly-detection-worker-cli \
--run-once \
--input /tmp/group.json \
--output /tmp/anomalies.json
--run-once raises a usage error unless both --input and --output are given. If the input file
contains inline sessions, no database or embedding model is needed.
Build a wheel:
uv build workers/anomaly-detection-worker
Configuration
Copy .env.template to .env, or pass flags.
| Environment variable | CLI flag | Description |
|---|---|---|
RABBITMQ_URL |
--rabbitmq-url |
Full connection URL |
ANOMALY_INPUT_QUEUE |
--input-queue |
Input queue name |
ANOMALY_FEEDBACK_QUEUE |
--feedback-queue |
Optional feedback queue |
EMBEDDING_MODEL |
--embedding-model |
Needed whenever sessions is empty |
ANOMALY_CONFIG_FILE |
--config-file |
YAML layer config |
| — | --max-sessions |
Stop after N messages (-1 = unlimited) |
NEO4J_HOST, NEO4J_PORT, NEO4J_USERNAME, NEO4J_PASSWORD, NEO4J_DB |
(API-managed) | Read by the API DAL, not by CLI flags |
The YAML config selects which layers run and which model each uses. Omitting it is equivalent to:
text:
model_name: isolation_forest
graph:
model_name: isolation_forest
metric:
model_name: elliptic_envelop
To run only the metric layer with a different model:
metric:
model_name: local_outlier_factor
Methodology
Source: dem/src/dem/anomaly/.
Each layer turns its sessions into a numeric feature matrix, hands it to a scikit-learn outlier
model, and reads the +1 / -1 labels back:
sessions in the group
│
├─ 1. featurise ──► matrix X of shape (n_sessions, n_features)
│
├─ 2. fit_predict ──► label per session: +1 inlier, -1 outlier
│
└─ 3. map labels back to session ids
fit_predict is used rather than fit then predict, meaning the model is fitted on the very
data it judges. There is no training phase and no persisted model; each group is scored in
isolation.
The models
AnomalyDetector (base.py) exposes three, keyed by the model_name config value. No extra
keyword arguments are forwarded, so scikit-learn's defaults apply throughout.
isolation_forest — default for text and graph
Builds 100 random trees. Each tree recursively splits the data by picking a random feature and a
random split value between that feature's min and max, until points are isolated. A point's
anomaly score is derived from E[h(x)], its average isolation depth across the forest, normalised
against c(n), the expected depth in a binary search tree of the same size:
s(x) = 2 ^ ( - E[h(x)] / c(n) )
Shallow depth gives a score near 1, so points that get isolated after only a few splits sit in
sparse regions and score high. With
contamination="auto" the decision offset is fixed at -0.5, which corresponds to s(x) ≈ 0.5 —
the score of a point at the expected average depth. The number of outliers is therefore not fixed
in advance. max_samples="auto" caps each tree's training subsample at 256 points.
Isolation Forest is the default for the two high-dimensional layers because its cost grows linearly with dimensionality and it makes no distributional assumption.
elliptic_envelop — default for metrics
Fits a robust Gaussian by Minimum Covariance Determinant: it searches for the subset of points whose covariance matrix has the smallest determinant (the tightest ellipsoid), and estimates the mean and covariance from that subset only, so a few extreme values cannot drag the fit. Points are then scored by squared Mahalanobis distance:
d(x) = (x - mu)^T * Sigma^-1 * (x - mu)
The cutoff comes from contamination, which defaults to 0.1 in scikit-learn.
Warning: Elliptic Envelope flags roughly 10% of sessions by construction Because
contaminationdefaults to0.1and no override is passed, this model marks about a tenth of each group as outliers for every metric, whether or not any session is genuinely unusual. Its own documentation also notes it is "not expected to yield meaningful results whenn_samples > n_features ^ 2", and the metric layer has exactly one feature. Small groups with a genuinely tight metric spread are the worst case here; the explicit guards below catch the degenerate ones, but not this general bias.
local_outlier_factor — opt-in
Compares the local density around a point to the density around its 20 nearest neighbours:
LOF(x) = mean over neighbours m of ( local_density(m) / local_density(x) )
LOF ≈ 1 means the point is as tightly packed as its neighbourhood; substantially greater than 1
means it sits in a comparatively sparse pocket. Unlike the other two this is a local criterion, so
it can flag a point that is unusual for its own cluster even when a group contains several distinct
clusters.
Text layer
TextualAnomalyDetector (textual_detector.py). The feature matrix is simply each session's
output_embedding stacked into an (n_sessions, embedding_dim) array — no dimensionality
reduction or rescaling. Sessions whose generated output is semantically unlike the rest of the
group come out as outliers.
Guards: if every embedding is identical, the model is skipped and all sessions are inliers.
Graph layer
GraphAnomalyDetector (graph_detector.py). Execution graphs are not numeric, so they are turned
into a bag of execution paths.
Step 1 — enumerate paths. _extract_path_ngrams runs a depth-first walk from every node,
emitting the label sequence at every step, up to max_path_length = 10 nodes. Labels come from
each node's agent_id (spaces replaced by _, missing values become unknown) and are joined
with __:
graph: A -> B -> C
paths: "A", "A__B", "A__B__C", "B", "B__C", "C"
All of a session's paths are then joined by spaces into one document string.
Step 2 — vectorise. CountVectorizer with token_pattern=r"(?u)\b\w[\w-]*\b" builds a
vocabulary of every path observed anywhere in the group and counts occurrences per session. The
custom token pattern keeps hyphenated agent identifiers intact instead of splitting them; __ is
made of word characters, so a whole path stays a single token.
Step 3 — score. The resulting count matrix goes into the model exactly like any other feature matrix.
The consequence of using paths rather than edges: a session is anomalous when it contains a sequence of steps nobody else took. Sharing every individual edge with the rest of the group but chaining them in a novel order is still enough to be flagged.
Guards: an empty vocabulary (all graphs empty) means all sessions are inliers; a single graph is reported as an inlier without fitting; identical count vectors across the group also short-circuit to all inliers.
Warning: Cyclic execution graphs are expensive The walk keeps no visited set, so it does not stop at cycles — it is bounded only by
max_path_length. A graph with a loop generates every repetition of that loop up to depth 10, and a densely connected graph generates a combinatorial number of paths. Sessions whose execution graphs contain retry loops are where this layer gets slow.
Metric layer
MetricAnomalyDetector (metric_detector.py). Each metric name is handled independently, as a
one-dimensional problem: values are collected across the group, reshaped to (n_sessions, 1), and
scored.
Before fitting, values that are non-numeric, NaN or infinite are dropped while their original
positions are remembered, so the surviving labels still map back to the right sessions.
Three explicit guards return "everything is an inlier" rather than a model result:
| Condition | Why |
|---|---|
| Fewer than 2 valid values | Nothing to compare against |
All values identical (n_unique == 1) |
Zero variance, so there is no notion of distance |
| Model raises "covariance matrix of the support data is equal to 0" or "Found array with 0 sample(s)" | Degenerate MCD fit on near-constant data |
The same covariance-zero fallback is implemented in the text detector.
Caveats
- No minimum group size. Unlike consistency and normal behaviour, this worker does not skip small groups. A group of one produces a report listing that single session as an inlier.
- Results are not stable across runs. Isolation Forest seeds its trees randomly and no
random_stateis set, so re-analysing the same group can shift the inlier/outlier split. - The idempotency pre-check only runs on the database path. Passing inline
sessionsbypassesanalysis_pre_check, so reports are recomputed and re-MERGEd on every call. - Zero reports still counts as success. When no layer produces a report the worker returns
Truewithout replacingoutput_messages, so run-once mode writes the input message back out unchanged rather than anAnomalyDetectionOutputMessage.