Skip to content

Embedding Worker

Standalone embedding worker for processing telemetry sessions with vector embeddings.

Source: workers/embedding-worker

Overview

The embedding worker computes, for each State in a session's trajectory, an embedded vector of its textual content, and writes the resulting embeddings back to the knowledge graph. These embeddings are what grouping-worker later uses to assign the session to a semantic group.

It sits right after norm-worker in the ingestion pipeline (see the architecture overview and pipelines docs), running in parallel with mce-worker:

  • Input queue: new_session_to_embedding
  • Output queue: new_session_to_grouping
  • CLI: embedding-worker-cli

Internally, EmbeddingWorker (workers/embedding-worker/src/embedding_worker/worker.py) subclasses the shared worker-base BaseWorker. On each message it fetches the session's state content from the knowledge graph (via oxp.client.dal), delegates to EmbeddingWrapper.process_session() (workers/embedding-worker/src/embedding_worker/wrapper/embedding_wrapper.py), and ingests the resulting vectors back through the same DAL. EmbeddingWrapper is a thin adapter around the dem.embedding submodule's SentenceTransformerEmbedder / OpenAIEmbedder classes: it builds the selected embedder, encodes each state's content (falling back to a placeholder string OXP_EMPTY_STATE for empty content), and attaches the model name and vector to each record before ingestion.

Usage

# Queue-based mode (long-running worker)
embedding-worker-cli \
  --rabbitmq-url amqp://<user>:<password>@rabbitmq:5672/ \
  --input-queue new_session_to_embedding \
  --output-queue new_session_to_grouping

# Run-once mode (Argo Workflows compatible)
embedding-worker-cli \
  --run-once \
  --input /tmp/input-message.json \
  --output /tmp/embedding-output-message.json

# Test liveness
embedding-worker-cli --test

Configuration

Copy workers/embedding-worker/.env.example to .env and fill in the values, or pass everything via CLI flags / environment variables:

Env var CLI flag Description
RABBITMQ_URL --rabbitmq-url Full RabbitMQ connection URL
EMBEDDING_INPUT_QUEUE --input-queue Input queue name (default new_session_to_embedding)
EMBEDDING_OUTPUT_QUEUE --output-queue Output queue(s), comma-separated (default new_session_to_grouping)
EMBEDDING_FEEDBACK_QUEUE --feedback-queue Optional feedback queue
EMBEDDING_EMBEDDER_TYPE --embedder-type SentenceTransformerEmbedder (default) or OpenAIEmbedder
EMBEDDING_MODEL --embedding-model Model identifier (default all-MiniLM-L6-v2; e.g. azure/text-embedding-3-small for OpenAIEmbedder)
NEO4J_URI / NEO4J_USERNAME / NEO4J_PASSWORD / NEO4J_DATABASE — Neo4j connection, used by the KG DAL
AI_GATEWAY_API_KEY / AI_GATEWAY_BASE_URL — Required when using OpenAIEmbedder
MAX_INFLIGHT_MESSAGES --max-inflight-messages Sessions processed concurrently per worker process (default 16; each in-flight session issues one batched embedding-API call)

Also supported: the common worker flags --max-sessions, --debug (also dumps embeddings to /tmp/embedding_embeddings_<session_id>.csv), --run-once, --input, --output, --test, shared by all workers.

Build & Run

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

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

Development

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