Worker Base
Shared RabbitMQ worker framework: every standalone worker in workers/ subclasses its BaseWorker.
Source: worker-base
Overview
worker-base is not a worker itself — it's the runtime every worker in the ingestion and
periodic pipelines (see the pipelines documentation) is built on
top of: async RabbitMQ consume/produce via aio_pika, Pydantic message schemas, feedback-queue
events, and an alternate single-message "Argo" mode for local testing or an Argo Workflows step.
All 11 workers under backend/workers/ (analysis-worker, anomaly-detection-worker,
consistency-worker, embedding-worker, grouping-worker, hierarchical-grouping-worker,
intelligence-worker, mce-worker, norm-worker, normal-behaviour-worker,
stateful-eval-worker) subclass BaseWorker and only implement handle_message() plus their own
CLI.
worker-base itself ships no CLI or argparse/click layer — the common flags every worker exposes
(--run-once, --input, --output, --test, --debug, --max-sessions,
--max-inflight-messages, --rabbitmq-url) are a convention each worker's own CLI module
implements (see e.g. workers/norm-worker/src/norm_worker/cli.py), built on top of the helper
functions worker-base provides in worker_base.utils (below).
Public API
BaseWorker (worker_base.base_worker)
Constructor (verified against base_worker.py):
BaseWorker(
rabbitmq_url: str,
input_queue: str,
output_queue: list[str] = [],
feedback_queue: str | None = None,
message_limit: int = -1, # -1 = unlimited
max_inflight_messages: int = 1,
output_message_ttl_seconds: float | None = None,
)
On construction it also tries to obtain a Neo4j connector via
oxp.dependencies.get_neo4j_connector() and wraps it in an OXPApiDALAdapter as
self.db_handler; if that fails (e.g. no DB configured), db_handler is left None and a
warning is logged instead of raising.
Key methods/attributes a subclass uses:
self.name— worker name used in logs and feedback messages.self.input_message_class— theBaseMessagesubclass (see below) used to deserialize incoming messages.async handle_message(msg) -> bool | None— must be overridden; the base implementation raisesNotImplementedError. ReturnTrueto propagateself.output_messagesto the output queue(s),Falsefor a successful no-op (nothing propagated),Noneto signal an error (aerrorfeedback event is emitted).self.emit_output_message(message)/self.output_messages— append to (or set) the message(s) to propagate; output messages are stored per-asynciotask (viacontextvars) so concurrent in-flight messages don't clobber each other's output.async run()— starts the long-running RabbitMQ consumer (run_worker()): declaresinput_queue, sets QoS/prefetch frommax_inflight_messages, and processes messages concurrently up to that limit, sendingstart/complete/errorevents tofeedback_queue(if set) around eachhandle_message()call.async run_argo_message(raw_input, output_path)— single-message (or batch-of-messages, ifraw_inputis a JSON array) mode with no RabbitMQ involved: validates the input againstinput_message_class, callshandle_message(), and writes the serialized output message(s) tooutput_path. This is what every worker's--run-once --input ... --output ...CLI flags drive.async propagate_message()/send_message_to_all()/send_message_to_queue()— fan the currentoutput_messagesout toself.output_queue(or a specific queue).
Messages (worker_base.queue_message)
Pydantic models, all subclassing BaseMessage (job_id, workflow_id, local_file, all
optional) with write_message() (publish to a channel/queue, persistent delivery, optional
ttl_seconds), from_body() (deserialize) and dump_message():
| Class | Adds |
|---|---|
BaseQueueMessage |
session_id (required) — the common input/output message for most workers |
BaseTriggerMessage |
application_id (required) — used by self-triggered periodic workers |
MasMetricTriggerMessage |
metric_name, mas_name (optional) |
ImpactAssessmentTrainingMessage |
end_to_end_mode, is_training_step |
MultiSessionQueueMessage |
session_ids (list), group_id |
SessionDetailMessage / SessionDetail |
metrics, input/output content and embeddings, execution graph |
SessionGroupMessage |
group_id, group_hash, sessions |
FeedbackMessage (a plain class, not a Pydantic model) carries session_id, event
(start/complete/error), workflow, workflow_id, queue_name and is what BaseWorker
publishes to feedback_queue around each handle_message() call.
DAL adapter (worker_base.dal_adapter.OXPApiDALAdapter)
A compatibility adapter that exposes the knowledge-graph operations workers need
(ingest_anomaly_report, ingest_consistency_report, ingest_normal_behaviour_report,
ingest_insights, get_analysis_data_for_semantic_group, get_semantic_group_hierarchy,
get_session_io_embeddings, get_mas_name, get_all_application_ids,
generate_insights_from_query, wait_for_hierarchical_grouping_unlock, analysis_pre_check),
delegating to the stateless functions in oxp.client.dal (from the api library)
against an injected Connector. Accessed as self.db_handler on a worker instance.
Utilities (worker_base.utils)
Helpers shared by every worker's CLI:
resolve_rabbitmq_url(explicit_url=None)— resolves--rabbitmq-urlin order: explicit value →RABBITMQ_URLenv var → splitRABBITMQ_HOST/RABBITMQ_PORT/RABBITMQ_USER/RABBITMQ_PASSWORD/RABBITMQ_VHOSTenv vars. Raises if nothing is configured.resolve_max_inflight_messages(cli_value, *, error_subject)—-1resolves to the available CPU count (get_available_cpus()); also readable from theMAX_INFLIGHT_MESSAGESenv var.resolve_neo4j_auth(...)— resolves Neo4j credentials fromNEO4J_AUTH(user/password) or the splitNEO4J_USERNAME/NEO4J_PASSWORDenv vars.configure_worker_logging(*, debug=False)— installs a coloredStreamHandlerformatter on the root logger and quietsneo4j.notifications.mask_url_password(url)— redacts the password in a connection URL for safe logging.
Development
cd worker-base
uv sync
uv run pytest tests/ -v
Tests live in worker-base/tests/ (test_base_worker.py, test_queue_message.py).