Skip to content

Hierarchical Grouping Worker

Periodic worker that (re)computes the hierarchical semantic-group structure over all of a MAS application's sessions, and flags which groups need (re-)analysis downstream.

Source: workers/hierarchical-grouping-worker

Overview

Unlike the per-session grouping-worker, which assigns one freshly-ingested session to an existing group as part of the ingestion pipeline, the hierarchical grouping worker runs periodically over all sessions of an application to build and maintain the semantic-group hierarchy itself — creating new groups, merging/splitting existing ones, and re-arranging the tree as new sessions accumulate (see the architecture overview). Both workers share the same clustering logic, dem.grouping.SemanticGrouper/Hierarchy (see the analysis component); the hierarchical worker additionally uses an LLM to name and summarize groups.

It is a self-triggering worker (see the pipelines documentation): HierarchicalGroupingWorker.run() (src/hierarchical_grouping_worker/worker.py) starts a background _periodic_self_trigger() task alongside the normal RabbitMQ consume loop. That task sleeps for periodic_trigger_interval_seconds (default 600s), then publishes one BaseTriggerMessage per known MAS application ID back onto its own input queue — there is no external scheduler. The output message TTL is automatically set to 1.5x the trigger interval.

  • Self-trigger / input queue: new_session_to_periodic_grouping
  • Output queue: new_session_to_analysis
  • CLI: hierarchical-grouping-worker-cli

On each trigger, HierarchicalGroupingWorker.handle_message() calls HierarchicalGroupingWrapper.compute_semantic_groups(application_id) to recompute the hierarchy, then merges that result with get_semantic_groups_needing_analysis() (groups whose hierarchy changed but haven't been re-analyzed yet) and emits one SessionGroupMessage per group — this is what feeds analysis-worker for anomaly detection, consistency, and normal-behaviour analysis. HierarchicalGroupingWorker subclasses the shared worker-base BaseWorker.

Usage

cp .env.template .env
# edit .env with your Neo4j and LLM settings

# RabbitMQ mode (long-running, self-triggering worker)
hierarchical-grouping-worker-cli \
  --input-queue new_session_to_periodic_grouping \
  --output-queue new_session_to_analysis \
  --embedding-model azure/text-embedding-3-small

# Test liveness
hierarchical-grouping-worker-cli --test

Configuration

Neo4j connection settings are resolved by the API layer (oxp.dependencies) from environment variables; the worker does not expose Neo4j CLI flags.

Env var CLI flag Default Description
RABBITMQ_URL --rabbitmq-url unset Full RabbitMQ connection URL
HIERARCHICAL_GROUPING_INPUT_QUEUE --input-queue new_session_to_periodic_grouping Self-trigger/input queue
HIERARCHICAL_GROUPING_OUTPUT_QUEUE --output-queue new_session_to_analysis Output queue(s), comma-separated
HIERARCHICAL_GROUPING_FEEDBACK_QUEUE --feedback-queue unset Optional feedback queue
PERIODIC_TRIGGER_INTERVAL_SECONDS --periodic-trigger-interval-seconds 600 How often the worker self-triggers a recompute; also sets output message TTL (1.5x)
HIERARCHICAL_GROUPING_EMBEDDING_MODEL --embedding-model azure/text-embedding-3-small Embedding model used for clustering
HIERARCHICAL_GROUPING_MAX_NEIGHBORS --max-neighbors 10 Max neighbours considered per session
HIERARCHICAL_GROUPING_MIN_SAMPLES --min-samples 20 Min samples required to form a cluster
HIERARCHICAL_GROUPING_MAX_DISTANCE --max-distance 0.15 Max cosine distance for grouping
LLM_BASE_MODEL_URL_MCE --llm-base-url unset LLM base URL, used to name/summarize groups
LLM_MODEL_NAME --llm-model-name gpt-4o LLM model identifier
OPENAI_API_KEY --llm-api-key unset LLM API key
MAX_INFLIGHT_MESSAGES --max-inflight-messages 1 Max app jobs processed concurrently in one worker process

Also supported: the common worker flags --run-once, --input, --output, --test, --debug, shared by all workers.

Build & Run

# Build a wheel (run from the monorepo root)
uv build workers/hierarchical-grouping-worker

# Sync dependencies and run the CLI directly
uv run --directory workers/hierarchical-grouping-worker hierarchical-grouping-worker-cli --test

Development

uv run --package hierarchical-grouping-worker --extra dev pytest workers/hierarchical-grouping-worker/tests/ -v