Intelligence Worker
Standalone periodic worker that executes intelligence catalog templates against the knowledge graph and ingests the resulting Insight nodes.
Source: workers/intelligence-worker
Overview
IntelligenceWorker (workers/intelligence-worker/src/intelligence_worker/worker.py) subclasses the shared worker-base BaseWorker. Unlike the ingestion-pipeline workers, it isn't chained to an upstream queue — it self-triggers: an internal asyncio task sleeps for periodic_trigger_interval_seconds (default 600s), then looks up every known application (MAS) id via db_handler.get_all_application_ids() and publishes one BaseTriggerMessage per application id back onto its own input queue.
Each trigger is handled by IntelligenceWrapper.generate_insights() (workers/intelligence-worker/src/intelligence_worker/wrapper/intelligence_wrapper.py), which wraps dem.intelligence:
- On startup,
load_templates_from_disk(catalog_root)(from the intelligence catalog library) loads everyInsightTemplateModel— each one anameTemplate/descriptionTemplate, akgQuery, and optionallabels/scope/priority/targetNodeId— from the catalog directory. - For the triggered application, every template's
kgQueryis run against the knowledge graph (db_handler.generate_insights_from_query), one row per match. - Each row's variables are substituted into the template's
$variableplaceholders to build anInsightnode (id = sha256 of template id + target node id + rendered name, for de-duplication). - The unique resulting
Insightnodes are ingested into the knowledge graph (db_handler.ingest_insights).
This is how SemanticGroup, AnomalyReport, NormalBehaviourReport, and Metric nodes (written earlier in the pipeline by grouping, analysis, and mce workers) get surfaced as rendered, human-readable insights — see the periodic operations section of the pipelines doc.
- Self-trigger queue:
new_session_to_intelligence - Output queue: none — results are persisted directly to the knowledge graph
- CLI:
intelligence-worker-cli - Config:
--catalog-root(path to the intelligence catalog)
Usage
# Queue-based mode (long-running worker)
intelligence-worker-cli \
--rabbitmq-url amqp://<user>:<password>@rabbitmq:5672/ \
--input-queue new_session_to_intelligence \
--catalog-root intelligence-catalog
# Run-once mode (Argo Workflows compatible)
intelligence-worker-cli \
--run-once \
--input /tmp/input.json \
--output /tmp/output.json \
--catalog-root intelligence-catalog
# Test liveness
intelligence-worker-cli --test
Configuration
| Env var | CLI flag | Description |
|---|---|---|
RABBITMQ_URL |
--rabbitmq-url |
Full RabbitMQ connection URL |
RABBITMQ_HOST / RABBITMQ_PORT |
(env only) | Split RabbitMQ connection settings, alternative to RABBITMQ_URL |
RABBITMQ_USER / RABBITMQ_PASSWORD |
(env only) | Required in split mode |
RABBITMQ_VHOST |
(env only) | Optional RabbitMQ virtual host |
INTELLIGENCE_INPUT_QUEUE |
--input-queue |
Self-trigger/input queue name (default new_session_to_intelligence) |
INTELLIGENCE_FEEDBACK_QUEUE |
--feedback-queue |
Optional feedback queue |
INTELLIGENCE_CATALOG_ROOT |
--catalog-root |
Root directory containing intelligence catalog templates (default .) |
MAX_INFLIGHT_MESSAGES |
--max-inflight-messages |
Max application jobs processed concurrently per worker process (default 1, -1 = auto CPU) |
PERIODIC_TRIGGER_INTERVAL_SECONDS |
--periodic-trigger-interval-seconds |
How often the worker self-triggers (default 600); also sizes the published message TTL (1.5x interval) |
| (none) | --template-max-concurrency |
Max insight templates executed concurrently per trigger (default 8) |
Also supported: the common worker flags --run-once, --input, --output, --test, --debug, and --max-sessions, shared by all workers.
Build & Run
# Build a wheel (run from the monorepo root)
uv build workers/intelligence-worker
# Sync dependencies and run the CLI directly
uv run --directory workers/intelligence-worker intelligence-worker-cli --test
Development
uv run --package intelligence-worker --extra dev pytest workers/intelligence-worker/tests/ -v