# Embedding Engine

Pipeline-stage service for the Adpilot indexing platform: consumes Mongo artifacts and Redis stream events, builds homogeneous `embedding_records`, generates OpenAI embeddings, and stores vectors in **Qdrant**.

Legacy HTTP routes still accept pre-chunked docs payloads from early boilerplate work; pipeline integration is the primary direction (see TW-78).

## Configuration

Copy `.env.example` to `.env`. Key groups:

| Group | Variables |
| --- | --- |
| App | `APP_NAME`, `ENVIRONMENT` (or `ENV`), `PORT`, `LOG_LEVEL` |
| MongoDB | `MONGO_URI`, `DATABASE_NAME` |
| OpenAI | `OPENAI_API_KEY`, `EMBEDDING_MODEL`, `EMBEDDING_PROVIDER`, `EMBEDDING_DIMENSION`, `EMBEDDING_BATCH_SIZE` |
| Qdrant | `QDRANT_URL`, `QDRANT_COLLECTION_NAME`, `QDRANT_API_KEY` |
| Retrieval | `RETRIEVE_DEFAULT_TOP_K`, `RETRIEVE_DEFAULT_SCORE_THRESHOLD`, `RETRIEVE_MAX_TOP_K`, `RETRIEVE_GRAPH_EXPAND_MAX_NODES`, `RETRIEVE_GRAPH_SCORE_DECAY` |
| Redis streams | `REDIS_URL`, `REDIS_STREAM_PREFIX`, `CONSUMER_GROUP` |
| Celery | `CELERY_BROKER_URL`, `CELERY_RESULT_BACKEND`, `CELERY_WORKER_CONCURRENCY` |
| Tokens / retry | `TOKENIZER_MODEL`, `MAX_TOKENS_PER_CHUNK`, `RETRY_MAX_ATTEMPTS` |

In `development` / `local`, missing secrets are allowed. In `production`, startup validation fails fast if required values are empty.

Contracts live in `app/models/embedding_record.py` and `app/models/stream_events.py` (see `indexing-pipeline.md` §4.5).

## Source Types (pipeline)

Normalized `source_type` on embedding records:

- `code`
- `docs`
- `commit`

Legacy HTTP API may still accept `repo_doc` from docs-ingestion boilerplate.

## Project Structure

- `app/core/config.py`: validated settings + Redis stream name helper
- `app/models/embedding_record.py`: pipeline embedding record contract
- `app/services/embedding_record_builder.py`: maps `doc_chunks`, `code_chunks`, `commit_analyses` → `embedding_records`
- `app/qdrant/`: Qdrant client, collection bootstrap, deterministic point IDs, vector upsert/delete
- `app/models/stream_events.py`: `graph.artifact.ready` / `commit.analysis.ready` events
- `app/api/routes/embeddings.py`: queue and status endpoints
- `app/services/embedding_service.py`: OpenAI provider with retries, token truncation, dimension validation
- `app/services/record_embedding_service.py`: batch embed records and index vectors in Qdrant
- `app/retrieval/`: read path — query embed, Qdrant search, Mongo hydration, grouped response
- `app/api/routes/internal_retrieve.py`: internal `POST /internal/retrieve` (RAG orchestrator only)
- `app/workers/celery_app.py`: Celery app + Redis broker
- `app/workers/embedding_tasks.py`: async embedding tasks with retry

## Requirements

Install dependencies:

```bash
pip install -r requirements.txt
pip install -r tests/requirements-test.txt
```

## Run Locally

Start Redis (and Mongo/Qdrant when wired in docker-compose):

```bash
redis-server
```

Start worker:

```bash
celery -A app.workers.celery_app worker --loglevel=info
```

Start API:

```bash
uvicorn app.main:app --reload --port 6004
```

Run unit tests:

```bash
pytest tests/unit -q
```

Run integration tests (in-memory pipeline E2E; no Docker required):

```bash
pytest tests/integration -m integration -q
```

Optional live checks against a running `docker compose` stack:

```bash
docker compose up -d --build
INTEGRATION_LIVE=1 pytest tests/integration -m live -q
```

## Docker

From the repo root (`adpilot-indexing-embedding-engine.com`):

```bash
cp .env.example .env
mkdir -p logs
docker compose up -d --build
```

Services: `app`, `celery_worker`, `stream_consumer`, `mongo`, `redis`, `qdrant`.

Check the API:

```bash
curl http://localhost:6004/healthz
curl http://localhost:6004/readyz
curl http://localhost:6004/api/v1/health/
curl http://localhost:6004/api/v1/health/readyz
curl http://localhost:6004/api/v1/version
```

Check dependencies:

```bash
docker compose exec mongo mongosh --eval "db.adminCommand('ping')"
docker compose exec redis redis-cli ping
curl http://localhost:6333/healthz
```

View logs:

```bash
docker compose logs -f app
docker compose logs -f celery_worker
docker compose logs -f stream_consumer
```

Stop:

```bash
docker compose down
```

## API Endpoints

### Health and meta

- `GET /healthz`: liveness probe
- `GET /readyz`: readiness probe (Mongo, Redis, Qdrant)
- `GET /api/v1/health/`: liveness under API prefix
- `GET /api/v1/health/readyz`: readiness under API prefix
- `GET /api/v1/readyz`: readiness alias under API prefix
- `GET /api/v1/version`: service name, version, environment

### Internal retrieval (TW-96 / TW-129)

Internal-only endpoint for the RAG orchestrator. Not exposed via API Gateway.

**Contract status:** Frozen 2026-06-10 ([TW-129](https://admedia-jira.atlassian.net/browse/TW-129)). Canonical types: `app/retrieval/models.py`. Full spec: `docs/retrieval-embedding-service-contracts.md` in the composer workspace (RAG TW-97 integrates against this version).

- `POST /internal/retrieve`: embed query text, search Qdrant, hydrate Mongo, return grouped context

Example:

```bash
curl -X POST http://localhost:6004/internal/retrieve \
  -H "Content-Type: application/json" \
  -d '{
    "repo_id": "AD/example-repo",
    "snapshot_id": "snap_abc123",
    "query_text": "How does configuration loading work?",
    "filters": { "source_types": ["code", "docs"] },
    "top_k": 20,
    "score_threshold": 0.25,
    "graph_expand": false,
    "request_id": "req_local_demo"
  }'
```

Response includes `hits`, `code_snippets`, `doc_excerpts`, `related_commits`, and `metadata` latency fields.

**Error codes:** `404 snapshot_not_found` (no `embedding_records` for repo/snapshot), `422 validation_error`, `502` embedding/qdrant, `503 mongo`. Indexed snapshot with zero hits above threshold returns `200` with empty arrays.

### Internal Lucos chunk indexing (TW-218)

Internal-only endpoints consumed by `lucos-cloud-client` during workspace indexing. Protected by `X-Internal-Secret` (`GATEWAY_INTERNAL_SERVICE_SECRET`).

- `POST /internal/lucos/chunks/embed`: embed Lucos workspace chunks into Mongo + Qdrant
- `POST /internal/lucos/chunks/delete`: delete vectors by deterministic `record_id`

Example:

```bash
curl -X POST http://localhost:6004/internal/lucos/chunks/embed \
  -H "Content-Type: application/json" \
  -H "X-Internal-Secret: $GATEWAY_INTERNAL_SERVICE_SECRET" \
  -d '{
    "chunks": [{
      "chunk_id": "src/auth/middleware.ts:10-42:abc123",
      "chunk_hash": "sha256:deadbeef",
      "repo_id": "lucos:ws:a1b2c3d4e5f67890",
      "file_path": "src/auth/middleware.ts",
      "content": "export function authMiddleware() { ... }",
      "source_type": "code",
      "start_line": 10,
      "end_line": 42,
      "language": "typescript",
      "symbol_name": "authMiddleware",
      "snapshot_id": "workspace",
      "commit_sha": "abc123"
    }]
  }'
```

Deterministic vector key: `record_id = lucos:{repo_id}:{chunk_id}`. Re-uploading the same `chunk_hash` is idempotent (no duplicate OpenAI/Qdrant work). Lucos retrieval uses the existing `POST /internal/retrieve` contract with `repo_id`, `snapshot_id`, and `filters.source_types` of `code` / `docs`.

### Pipeline (TW-89)

- `GET /api/v1/runs/{run_id}`: embedding run status
- `GET /api/v1/runs?repo_id=&snapshot_id=`: list runs for a snapshot
- `GET /api/v1/records/{record_id}`: single embedding record
- `GET /api/v1/records?repo_id=&snapshot_id=`: list records (optional `source_type`, `embed_status`)
- `GET /api/v1/records/count/summary?repo_id=&snapshot_id=`: record counts by status
- `POST /api/v1/pipeline/snapshots/embed`: manually queue snapshot embedding (202)
- `POST /api/v1/pipeline/commits/embed`: manually queue commit-analysis embedding (202)

### Legacy embedding helpers

- `POST /api/v1/embeddings/embed`: queue one chunk
- `POST /api/v1/embeddings/embed-batch`: queue many chunks
- `GET /api/v1/embeddings/embed/{task_id}`: check task status
- `POST /api/v1/embeddings/embed-sync`: sync embedding (debug/small payloads)
- `GET /api/v1/embeddings/health`: service configuration health

### Example: queue one chunk

```bash
curl -X POST http://localhost:6004/api/v1/embeddings/embed \
  -H "Content-Type: application/json" \
  -d '{
    "chunk_id": "chunk_1",
    "chunk_text": "Example text",
    "source_type": "repo_doc",
    "metadata": {"doc": "README.md"}
  }'
```

## Notes

- `/readyz` checks Mongo, Redis, and Qdrant connectivity (TW-87).
- A valid `OPENAI_API_KEY` is required for embedding generation in non-local environments.
- Retrieval logs include `request_id`, hit counts, and vector/hydrate/graph latency on each `/internal/retrieve` call.
- `graph_expand: true` traverses `snapshot_graphs` one hop to add linked code/doc/commit context (P6 / TW-102).
- Large merged graphs may be stored as a manifest in `snapshot_graphs` plus rows in `snapshot_graph_parts`; `snapshot_graph_repo.get_merged()` reassembles them transparently before traversal.
- When `filters.source_types` includes `commit`, commit vector scores receive `RETRIEVE_COMMIT_SCORE_BOOST`.

