Compare commits

..

3 commits
main ... main

Author SHA1 Message Date
itexpert228
408c7df005
Merge upstream main and resolve conflicts 2026-04-18 13:35:23 +03:00
itexpert228
70f647f17d
index commit 1 2026-04-18 12:38:29 +03:00
itexpert228
b9d9f38f2b
ignore DS_Store 2026-04-18 11:53:28 +03:00
37 changed files with 1651 additions and 2000 deletions

125
.ai_update/changes.md Normal file
View file

@ -0,0 +1,125 @@
# AI Change Log
Дата создания: 2026-04-18
## Как пользоваться
Этот файл - рабочий журнал изменений для Codex.
Перед новыми правками нужно прочитать этот файл и учитывать:
- что уже было изменено;
- какие файлы трогались;
- какие проверки запускались;
- какие ограничения и договоренности есть по задаче.
После каждой осмысленной правки нужно добавлять новую запись с:
- кратким описанием изменения;
- списком измененных файлов;
- результатом проверок;
- открытыми рисками или TODO, если они есть.
## Договоренности
- Не менять контракты `POST /index`, `POST /sparse_embedding`, `POST /search`.
- По текущей задаче фокус держать на `index`, если пользователь не просит иначе.
- Не трогать `docker-compose.yml` и инфраструктуру без отдельной просьбы.
- Не откатывать чужие или пользовательские изменения.
## Записи
### 2026-04-18 - разрешены конфликты с upstream/main
Что сделано:
- Объединены записи журнала из ветки PR и `upstream/main`.
- В `search/main.py` сохранены расширенные retrieval/rerank доработки из обеих веток.
- Убраны конфликтные маркеры после merge `upstream/main`.
- Сохранены безопасные лимиты retrieval и защита от падения rerank.
Измененные файлы:
- `.ai_update/changes.md`
- `search/main.py`
Проверки:
- `python3 -m py_compile index/main.py search/main.py`
- Прямой smoke `build_chunks` на `data/Go Nova.json`: 15 чанков, покрыто 25 из 25 сообщений.
- Pure smoke для `search/main.py` helper-функций через stub-модули.
Открытые риски:
- Полный интеграционный прогон с настоящими Qdrant/dense/rerank зависит от внешних сервисов.
### 2026-04-18 - index P1 and search retrieval/rerank
Что сделано:
- В `index/main.py` заменен символьный chunking на сборку чанков окнами сообщений.
- Добавлен учет временного разрыва между сообщениями через `INDEX_TIME_GAP_SECONDS`.
- Overlap теперь работает по границам сообщений через `INDEX_CHUNK_OVERLAP_MESSAGES`, а не по хвосту строки.
- Сообщения рендерятся структурно: `author`, `time`, `thread`, `mentions`, флаги, `quote`, `forward`, `file`, `system_event`.
- `file_snippets` парсятся как JSON; в индекс попадают имя файла, mime, url, владелец и дата создания.
- `member_event` превращается в индексируемый системный текст.
- `page_content`, `dense_content`, `sparse_content` разведены.
- В `index/Dockerfile` старый `CHUNK_SIZE=10` заменен на реальные `INDEX_*` настройки chunking.
- В `search/main.py` исправлен runtime-баг с неинициализированным `must_conditions`.
- Основной query в search теперь берется из `question.search_text` с fallback на `question.text`.
- `question.variants`, `question.hyde`, `question.keywords` и entities используются как дополнительные dense/sparse запросы.
- Retrieval делает несколько prefetch-запросов и fusion через Qdrant.
- Rerank больше не выбрасывает retrieval-хвост.
- Retrieval points дедуплицируются по Qdrant point id.
- Финальные `message_ids` агрегируются по score, дедуплицируются и ограничиваются `top-50`.
Измененные файлы:
- `index/main.py`
- `index/Dockerfile`
- `search/main.py`
- `.ai_update/changes.md`
Проверки:
- `python3 -m py_compile index/main.py search/main.py`
- Прямой smoke `build_chunks` на `data/Go Nova.json`.
- Прямой smoke endpoint-функции `index(...)` на `data/Go Nova.json`.
- Pure smoke для `search` helper-функций через stub-модули, потому в host env нет `qdrant_client` и `httpx`.
Результаты проверки индекса:
- Было 29 чанков, стало 15.
- Покрытие сообщений на `data/Go Nova.json`: 25 из 25.
- Системное сообщение с `member_event` больше не выпадает.
- `file_snippets` с `IMG_8471.webp` попадает в `sparse_content`.
- Quote и forward маркеры попадают в `dense_content`.
- `page_content`, `dense_content`, `sparse_content` больше не одинаковые.
Открытые риски:
- Полный интеграционный прогон `search` с настоящими Qdrant/dense/rerank локально не выполнялся.
- `date_range` фильтр включается только если установленный `qdrant_client` поддерживает `models.DatetimeRange`.
- Полный docker build локально не запускался.
### 2026-04-18 - создан журнал изменений
Что сделано:
- Создан файл `.ai_update/changes.md`.
- Зафиксировано, что до этого код не менялся, была только разведка репозитория и ТЗ.
Контекст по текущему состоянию:
- `index/main.py` тогда использовал символьный chunking.
- `render_message` брал только `message.text` и `parts[*].text`.
- `page_content`, `dense_content`, `sparse_content` были одинаковые.
- В примере `data/Go Nova.json` старый `build_chunks` покрывал 24 из 25 сообщений; системное сообщение с `member_event` выпадало из индекса.
Измененные файлы:
- `.ai_update/changes.md`
Проверки:
- Код не запускался, потому что создан только журнал.

View file

@ -1,13 +0,0 @@
# Local docker compose configuration
QDRANT_URL=http://qdrant:6333
QDRANT_COLLECTION_NAME=evaluation
QDRANT_DENSE_VECTOR_NAME=dense
QDRANT_SPARSE_VECTOR_NAME=sparse
EMBEDDINGS_DENSE_URL=http://83.166.249.64:18001/embeddings
RERANKER_URL=http://83.166.249.64:18001/score
# Fill either API_KEY or both OPEN_API_LOGIN and OPEN_API_PASSWORD.
API_KEY=
OPEN_API_LOGIN=
OPEN_API_PASSWORD=

5
.gitignore vendored
View file

@ -149,10 +149,6 @@ activemq-data/
# Environments
.env
.env.local
.env.*.local
!.env.example
.ai_update/
.envrc
.venv
env/
@ -218,3 +214,4 @@ __marimo__/
# Streamlit
.streamlit/secrets.toml
.DS_Store

View file

@ -1,44 +0,0 @@
# Repository Instructions
This repository uses a shared Codex collaboration workflow.
## Mandatory behavior for every Codex session
- Before doing any work, restore context from:
1. `README.md` and project docs if present
2. `.ai_update/current_status.md`
3. `.ai_update/handoff.md`
4. latest files in `.ai_update/sessions/`
5. `.ai_update/changelog.md`
6. git state (`git status`, recent `git log`)
- For any coding task in this repo, use the skill:
`team-sync-hackathon`
- After any meaningful change, update `.ai_update/` so another human or Codex session can continue with zero guesswork.
- Prefer continuing existing architecture and conventions over rewriting working code.
- Do not delete or reset teammates' work unless explicitly requested.
## Shared memory rules
Codex must treat `.ai_update/` as the canonical shared handoff area between:
- humans
- current Codex session
- future Codex sessions
If `.ai_update/` is missing or incomplete, create/fill it before large changes.
## Local startup rule
If local validation is needed and the project is not running:
- first try `./bin/codex-start`
- otherwise follow the startup instructions from the skill
## Safety rules for repo work
- Never change public API contracts unless explicitly requested.
- Never run destructive cleanup commands without explicit instruction.
- Never force-push without explicit instruction.
- Prefer small verifiable steps.

View file

@ -374,13 +374,14 @@ score = recall_avg * 0.8 + ndcg_avg * 0.2
Для локального запуска используйте `docker compose`.
Сначала подготовьте локальный env:
Перед запуском укажите учетные данные для внешнего dense/rerank API:
```bash
cp .env.example .env
export OPEN_API_LOGIN=...
export OPEN_API_PASSWORD=...
```
После этого заполните в `.env` либо `API_KEY`, либо пару `OPEN_API_LOGIN` / `OPEN_API_PASSWORD`.
Если эти переменные не заданы, `docker compose up` завершится с ошибкой.
Запуск:

290
doc/ai_update.md Normal file
View file

@ -0,0 +1,290 @@
# AI Update по ТЗ
Дата: 2026-04-18
## Что просмотрено
- `doc/ТЗа_хакатон_Индексация_и_поиск_по_сообщениям.pdf`
- `README.md`
- `docker-compose.yml`
- `index/main.py`, `search/main.py`
- `index/Dockerfile`, `search/Dockerfile`
- `index/Makefile`, `search/Makefile`
- `data/Go Nova.json`
По примеру данных:
- всего сообщений: `25`
- с `parts`: `14`
- с цитатами: `5`
- с пересланными сообщениями: `2`
- с `mentions`: `4`
- с `file_snippets`: `1`
- системных сообщений: `1`
Это важно, потому что в текущем коде часть этих сигналов либо не используется вообще, либо теряет смысл при индексации.
## Что уже соответствует ТЗ
1. В репозитории есть оба требуемых сервиса: `index` и `search`.
2. Обязательные endpoints реализованы:
- `GET /health`, `POST /index`, `POST /sparse_embedding` в `index/main.py`
- `GET /health`, `POST /search` в `search/main.py`
3. Контракты request/response по основным endpoint'ам не менялись и в целом совпадают с шаблоном и ТЗ.
4. `search` использует `Qdrant`, dense endpoint и reranker через HTTP.
5. В обоих Dockerfile sparse-модель предзагружается внутрь образа, что соответствует оффлайн-ограничению контейнеров.
6. `HOST` и `PORT` читаются из env, как требует ТЗ.
## Что отсутствует или реализовано частично
### 1. Обогащенный вопрос из ТЗ почти не используется
В `search/main.py:90-100` описаны поля:
- `search_text`
- `variants`
- `hyde`
- `keywords`
- `entities`
- `date_mentions`
- `date_range`
- `asker`
Но в реальном поиске используется только `question.text`:
- `search/main.py:307-316`
Это главный недобор относительно ТЗ. Само ТЗ явно дает эти поля как сигналы для retrieval, а код их сейчас просто игнорирует.
### 2. Метаданные чанков объявлены, но не участвуют в поиске
В `README.md:52-58` отдельно сказано, что в metadata чанка сохраняются:
- `participants`
- `mentions`
- `contains_forward`
- `contains_quote`
В `search/main.py:132-145` есть модель `ChunkMetadata`, но дальше она никак не используется в `query_points`. Поиск не делает:
- фильтрацию по `mentions`
- фильтрацию по `participants`
- учет `thread_sn`
- учет временного диапазона через `start`/`end`
- отдельную обработку quote/forward чанков
То есть сильный канал улучшения качества уже предусмотрен схемой, но сейчас не задействован.
### 3. Реранк отбрасывает часть кандидатов
Сейчас:
- retrieval берет до `20` чанков: `search/main.py:174-177`
- rerank берет только первые `10`: `search/main.py:278-296`
- после rerank возвращаются только эти `10`, а хвост `11-20` теряется: `search/main.py:321-328`
Это не нарушение контракта, но это реальная потеря recall.
### 4. Выдача не дедуплицируется и не ограничивается по полезному top-K
Сейчас `message_ids` просто конкатенируются:
- `search/main.py:323-328`
Проблемы:
- дубликаты message id не удаляются
- результаты не агрегируются по лучшему score сообщения
- нет явного ограничения на топ полезных `50`, хотя именно `K=50` участвует в метрике из ТЗ
Если один и тот же `message_id` попал в несколько чанков, он тратит место в выдаче.
### 5. Индексация пока очень базовая: фиксированные символьные чанки
В `index/main.py:118-190` чанки строятся просто по длине строки:
- `CHUNK_SIZE = 512`
- `OVERLAP_SIZE = 256`
- разбиение идет по символам, а не по сообщениям, тайм-гепам, тредам или смысловым блокам
Из-за этого:
- длинные пересланные сообщения и цитаты могут резаться в неудобных местах
- один и тот же смысловой блок может быть разнесен по чанкам неестественно
- overlap строится по хвосту текста, а не по границе сообщений
### 6. `page_content`, `dense_content` и `sparse_content` сейчас одинаковые
См. `index/main.py:180-186`.
ТЗ прямо оставляет это место как точку оптимизации качества, но пока этот резерв не используется.
### 7. Индексация берет только `text` и `parts[*].text`, остальное почти теряется
См. `index/main.py:99-115`.
Сейчас не используются как поисковые сигналы:
- `sender_id`
- `mentions`
- `file_snippets`
- `member_event`
- `thread_sn`
- `is_hidden`
- `is_system`
- явное различение `quote` и `forward`
Особенно важные пробелы:
- `member_event` у системных сообщений сейчас фактически пропадает, если обычного текста нет
- `file_snippets` не разбирается, хотя там могут быть имена файлов, URL и служебные поля
- запросы вида "кто писал", "кого упоминали", "какой файл/документ кидали" сейчас поддержаны слабо
### 8. Смысл `quote` и `forward` не маркируется
В `index/main.py:105-113` текст из `parts` просто подшивается в общий текст без явных маркеров вида:
- "цитата:"
- "пересланное сообщение:"
- "автор цитаты:"
В итоге dense/sparse видят просто общий текстовый комок. Для поиска по обсуждениям это ощутимая потеря контекста.
### 9. Есть расхождение между локальной инфраструктурой и ТЗ
По ТЗ для `search` ожидается `API_KEY`.
В коде это поддержано:
- `search/main.py:21-30`
- `search/main.py:42-47`
- `search/Makefile:10-15`
Но локальный `docker-compose.yml:57-60` требует `OPEN_API_LOGIN` и `OPEN_API_PASSWORD`.
Итог:
- сам сервис гибче ТЗ
- локальный compose не повторяет боевую схему из ТЗ один в один
Это не ломает контракт, но может запутать при локальной отладке.
### 10. Есть еще одна инфраструктурная несостыковка со сдачей
В `doc/upload_to_docker.md:45-48` явно сказано собирать образы с `--platform linux/amd64`.
Но `index/Makefile:16-18` и `search/Makefile:26-28` собирают без `--platform linux/amd64`.
На x86 это может пройти незаметно, а на ARM-машине дать неправильный образ для отправки.
### 11. Есть неоднозначность между README и PDF по sparse в `search`
- `README.md:211` говорит, что sparse-модель для `search` должна быть локально внутри образа
- PDF в формулировке требований к `Search Service` делает акцент, что обращения к dense/sparse/rerank идут через проверяющую систему
Текущий код следует логике README/example: sparse считается локально в `search/main.py:148-151` и `search/main.py:198-207`.
Я бы это не считал блокером, но как минимум это место стоит держать в голове как неоднозначное требование.
## Что улучшать в первую очередь
### Приоритет 1. Начать использовать все поля `question`
Минимально стоит задействовать:
- `search_text` как основной нормализованный запрос
- `variants` как дополнительные формулировки
- `hyde` как дополнительные dense-запросы
- `keywords` как основу для sparse
- `entities` для фильтров и lexical boost
- `date_range` и `date_mentions` для ограничения по времени
- `asker` как сигнал по людям и email
Самый логичный путь без смены стэка: несколько dense/sparse запросов + fusion в `Qdrant`.
### Приоритет 2. Перестроить chunking под структуру чата, а не под символы
Нужны чанки по:
- окнам сообщений
- временным разрывам
- границам thread/forward/quote
- ограничению на размер по сообщениям, а не только по символам
Для чатов это обычно дает больше пользы, чем любые косметические тюнинги rerank.
### Приоритет 3. Развести `page_content`, `dense_content`, `sparse_content`
Хорошая схема:
- `page_content`: читабельный исходный текст чанка
- `dense_content`: нормализованный текст с ролями, автором, маркерами quote/forward
- `sparse_content`: keyword-heavy версия с email, mentions, именами файлов, ссылками, документами, леммами
Сейчас эта возможность не используется вообще.
### Приоритет 4. Нормально собирать финальную выдачу
Нужно:
- не терять кандидатов после rerank
- удалять дубликаты `message_id`
- агрегировать по лучшему score сообщения или чанка
- отдавать осмысленный top-50
Это прямой выигрыш по Recall@50 и nDCG@50.
### Приоритет 5. Начать использовать metadata в `Qdrant`
Особенно полезно для:
- `mentions`
- `participants`
- `contains_quote`
- `contains_forward`
- `thread_sn`
- `start` / `end`
Для многих вопросов это позволит не просто "лучше ранжировать", а сразу отрезать нерелевантный шум.
### Приоритет 6. Превратить скрытые сигналы в индексируемый текст
Стоит отдельно материализовать:
- `member_event` в текст вида "пользователь X добавил Y"
- `file_snippets` в текст вида "файл: NAME, url: ..."
- автора сообщения
- список упомянутых пользователей
Сейчас эти сигналы либо не попадают в индекс, либо попадают слишком слабо.
## Что можно добавить по Python-библиотекам, не меняя стек
Стек `Qdrant` менять не нужно. Самые полезные добавки я бы смотрел такие:
- `pymorphy3` для лемматизации русских слов при подготовке `sparse_content`
- `razdel` для аккуратной токенизации русского текста
- `rapidfuzz` для точного lexical match по именам, email, документам, ссылкам и названиям
- `python-dateutil` или `dateparser` для нормализации дат, если захотите усиливать работу с `date_mentions`
- `tenacity` для аккуратных retry/timeout-оберток вокруг dense/rerank HTTP вызовов
Что не нужно делать:
- менять `Qdrant`
- тащить внешние LLM/API
- усложнять архитектуру ради "модности", пока не использованы базовые сигналы из самого ТЗ
## Короткий вывод
Сейчас репозиторий соответствует ТЗ как рабочий базовый шаблон, но почти не использует те сигналы, ради которых это ТЗ вообще интересно:
- обогащение вопроса
- metadata чанков
- структуру chat messages
- сигналы автора, упоминаний, файлов, системных событий, цитат и пересылок
Самый большой потенциал улучшения здесь не в замене базы или модели, а в трех вещах:
1. умный chunking
2. multi-query hybrid retrieval
3. использование metadata и нормальной сборки финального top-50

49
doc/todo.md Normal file
View file

@ -0,0 +1,49 @@
# TODO
## `search/main.py`
- [ ] P0: Переключить основной query на `question.search_text` с fallback на `question.text`
- [ ] P0: Подключить `question.variants` как дополнительные query-варианты
- [ ] P0: Подключить `question.hyde` как дополнительные dense-запросы
- [ ] P0: Подключить `question.keywords` как основу для sparse-запросов
- [ ] P0: Перестать терять retrieval-кандидатов после rerank
- [ ] P0: Дедуплицировать `message_ids` перед ответом
- [ ] P0: Ограничить финальную выдачу до `top-50`
- [ ] P0: Агрегировать score по `message_id`
- [ ] P2: Использовать `entities.people` и `entities.emails` для boost или фильтрации
- [ ] P2: Использовать `entities.documents`, `entities.names`, `entities.links` для lexical boost
- [ ] P2: Использовать `date_range` для фильтрации по `metadata.start` и `metadata.end`
- [ ] P2: Использовать `contains_quote` и `contains_forward` как сигналы ранжирования
- [ ] P2: Добавить multi-query fusion в `Qdrant`
- [ ] P2: Подобрать `prefetch`, `retrieve_k`, `rerank_limit`
- [ ] P3: Добавить retry и timeout политику для dense/rerank HTTP вызовов
## `index/main.py`
- [ ] P1: Перейти с символьного chunking на chunking по сообщениям
- [ ] P1: Учитывать time gap при сборке чанков
- [ ] P1: Маркировать в тексте `quote`, `forward`, автора сообщения и автора цитаты
- [ ] P1: Развести `page_content`, `dense_content`, `sparse_content`
- [ ] P1: Материализовать `sender_id` и `mentions` в индексируемый текст
- [ ] P1: Разбирать `file_snippets` и вытаскивать имя файла, mime и url
- [ ] P1: Разбирать `member_event` и превращать его в индексируемый текст
## `docker-compose.yml`
- [ ] P3: Привести локальный `docker-compose.yml` к схеме с `API_KEY`
- [ ] P3: Добавить `--platform linux/amd64` в сборку образов
## `doc/` (новый файл с регрессионными вопросами)
- [ ] P3: Зафиксировать набор локальных тестовых вопросов для проверки регрессий
## `search/requirements.txt`
- [ ] P4: Добавить `python-dateutil` или `dateparser`
- [ ] P4: Добавить `tenacity`
## `index/requirements.txt` и/или `search/requirements.txt`
- [ ] P4: Добавить `razdel`
- [ ] P4: Добавить `pymorphy3`
- [ ] P4: Добавить `rapidfuzz`

223
doc/todo_and_pipeline.md Normal file
View file

@ -0,0 +1,223 @@
# To-Do и целевой pipeline
Дата: 2026-04-18
## Цель
Поднять `Recall@50` и `nDCG@50` без смены стэка:
- оставить `Qdrant`
- оставить внешний dense endpoint
- оставить внешний reranker
- усиливать только индекс, retrieval, rerank и post-processing
## Приоритетный to-do
### P0. Быстрые и самые окупаемые правки
- [ ] Переключить основной запрос в `search` на `question.search_text` с fallback на `question.text`
- [ ] Подключить `question.variants` как дополнительные query-формулировки
- [ ] Подключить `question.hyde` как дополнительные dense-запросы
- [ ] Подключить `question.keywords` как основу для sparse-запроса
- [ ] Перестать терять кандидатов после rerank: возвращать не только top-10 rerank, но и хвост retrieval
- [ ] Дедуплицировать `message_ids` перед ответом
- [ ] Ограничить финальную выдачу осмысленным `top-50`
- [ ] Агрегировать score по `message_id`, а не просто конкатенировать ids из чанков
### P1. Улучшение индексации
- [ ] Перейти с символьного chunking на chunking по сообщениям
- [ ] Учитывать временные разрывы между сообщениями при сборке чанка
- [ ] Не смешивать в одном чанке слишком далекие по смыслу блоки
- [ ] Отдельно маркировать `quote`, `forward`, автора сообщения и автора цитаты
- [ ] Развести `page_content`, `dense_content`, `sparse_content`
- [ ] Материализовать `mentions` в текст и metadata
- [ ] Материализовать `sender_id` в индексируемый текст
- [ ] Разбирать `file_snippets` и вытаскивать имя файла, mime, url
- [ ] Разбирать `member_event` и превращать его в индексируемый текст
### P2. Улучшение retrieval и фильтрации
- [ ] Использовать `entities.people` и `entities.emails` для boost или фильтрации по `participants` и `mentions`
- [ ] Использовать `entities.documents`, `entities.names`, `entities.links` для lexical boost
- [ ] Использовать `date_range` для фильтрации по `metadata.start` и `metadata.end`
- [ ] Использовать `contains_quote` и `contains_forward` как дополнительные сигналы ранжирования
- [ ] Добавить multi-query fusion в `Qdrant` для dense и sparse запросов
- [ ] Подобрать новые значения `prefetch`, `retrieve_k`, `rerank_limit`
### P3. Инфраструктура и надежность
- [ ] Привести локальный `docker-compose.yml` к схеме с `API_KEY`, чтобы локальный запуск был ближе к ТЗ
- [ ] Добавить `--platform linux/amd64` в сборку образов
- [ ] Добавить retry и timeout политику для dense/rerank HTTP вызовов
- [ ] Зафиксировать набор локальных тестовых запросов для регрессии качества
### P4. Библиотеки, которые можно добавить без смены стэка
- [ ] `razdel` для токенизации русского текста
- [ ] `pymorphy3` для лемматизации при подготовке `sparse_content`
- [ ] `rapidfuzz` для точного match по именам, email, документам и ссылкам
- [ ] `python-dateutil` или `dateparser` для нормализации дат
- [ ] `tenacity` для retry вокруг внешних HTTP запросов
## Целевой pipeline индексации
### 1. Подготовка сообщения
На входе каждое сообщение должно раскладываться на сигналы:
- основной текст сообщения
- `parts[*].text`
- тип части: `text`, `quote`, `forward`
- `sender_id`
- `mentions`
- `file_snippets`
- `member_event`
- `thread_sn`
- флаги `is_system`, `is_quote`, `is_forward`
### 2. Нормализация и разметка
Перед chunking сообщение стоит приводить к структурированному виду, например:
- `author: ...`
- `mentions: ...`
- `quote: ...`
- `forwarded: ...`
- `file: ...`
- `system_event: ...`
Смысл не в красивом выводе, а в том, чтобы dense и sparse видели роль каждого куска текста.
### 3. Chunking
Целевой принцип:
- базовая единица не символ, а сообщение
- чанк собирается как окно из нескольких соседних сообщений
- окно режется по лимиту размера
- окно закрывается на большом time gap
- `forward` и длинные `quote` не должны ломать соседний контекст
- overlap должен работать по границам сообщений, а не по хвосту строки
### 4. Формирование трех видов текста
`page_content`:
- человекочитаемый текст чанка для payload
`dense_content`:
- нормализованный текст с автором, role-маркерами, quote/forward маркерами
`sparse_content`:
- keyword-heavy текст
- леммы
- email
- mentions
- имена файлов
- ссылки
- названия документов и сервисов
### 5. Metadata для Qdrant
В metadata стоит стабильно сохранять:
- `message_ids`
- `participants`
- `mentions`
- `thread_sn`
- `start`
- `end`
- `contains_quote`
- `contains_forward`
- `chat_id`
- `chat_type`
## Целевой pipeline поиска
### 1. Подготовка query
Собирать query не из одного поля, а из набора:
- основной запрос: `search_text` или `text`
- дополнительные dense-query: `variants` и `hyde`
- дополнительные sparse-query: `keywords`
- entity-сигналы: `people`, `emails`, `documents`, `names`, `links`
- time constraints: `date_range`, `date_mentions`
### 2. Query builder
Нужно строить несколько представлений запроса:
- dense-query для смысла
- sparse-query для точных слов и терминов
- filter/boost по metadata
### 3. Retrieval в Qdrant
Практическая схема:
1. Выполнить несколько dense prefetch по разным вариантам запроса
2. Выполнить несколько sparse prefetch по keyword-heavy запросам
3. Добавить filters по датам, mentions, participants, если это явно следует из вопроса
4. Объединить результаты через fusion
5. Забрать расширенный пул кандидатов для rerank
### 4. Rerank
Rerank должен работать не на слишком маленьком пуле. Целевой принцип:
- retrieval дает расширенный пул
- rerank сортирует top-N кандидатов
- хвост retrieval не теряется полностью
### 5. Агрегация к `message_id`
После rerank:
- собрать `message_ids` из чанков
- удалить дубликаты
- агрегировать лучший score на сообщение
- собрать финальный `top-50`
Это особенно важно, потому что метрики в ТЗ считаются именно по `message_id`, а не по chunk id.
## Порядок внедрения
### Этап 1. Quick wins
- использовать `search_text`, `variants`, `hyde`, `keywords`
- перестать терять кандидатов после rerank
- добавить dedup и top-50
### Этап 2. Пересборка индекса
- новый renderer сообщения
- новый chunking по сообщениям
- разные `page_content`, `dense_content`, `sparse_content`
### Этап 3. Metadata-aware retrieval
- filters по дате
- boost по mentions/participants
- учет `contains_quote` и `contains_forward`
### Этап 4. Тюнинг
- подобрать размеры чанков
- подобрать `retrieve_k`
- подобрать `rerank_limit`
- прогнать локальный набор контрольных вопросов
## Минимальный критерий готовности
Можно считать, что pipeline собран в рабочем виде, если:
- `search` использует не только `question.text`
- индексация не режет чанки посреди сообщения как основной механизм
- `message_ids` дедуплицируются
- финальная выдача ограничивается top-50
- retrieval умеет использовать хотя бы часть metadata
- локальная сборка и запуск не расходятся с ТЗ по критичным env и platform

64
doc/todo_people.md Normal file
View file

@ -0,0 +1,64 @@
Для того чтобы раскидать задачи между тремя людьми, я распределю их по сложности и приоритету.
### Человек 1 (Основной фокус на поиске):
#### `search/main.py`
* **P0**: Переключить основной query на `question.search_text` с fallback на `question.text`
* **P0**: Подключить `question.variants` как дополнительные query-варианты
* **P0**: Подключить `question.hyde` как дополнительные dense-запросы
* **P0**: Подключить `question.keywords` как основу для sparse-запросов
* **P0**: Перестать терять retrieval-кандидатов после rerank
* **P0**: Дедуплицировать `message_ids` перед ответом
* **P0**: Ограничить финальную выдачу до `top-50`
* **P0**: Агрегировать score по `message_id`
* **P2**: Использовать `entities.people` и `entities.emails` для boost или фильтрации
* **P2**: Использовать `entities.documents`, `entities.names`, `entities.links` для lexical boost
* **P2**: Использовать `date_range` для фильтрации по `metadata.start` и `metadata.end`
* **P2**: Использовать `contains_quote` и `contains_forward` как сигналы ранжирования
* **P2**: Добавить multi-query fusion в `Qdrant`
* **P2**: Подобрать `prefetch`, `retrieve_k`, `rerank_limit`
* **P3**: Добавить retry и timeout политику для dense/rerank HTTP вызовов
---
### Человек 2 (Основной фокус на индексации и разметке):
#### `index/main.py`
* **P1**: Перейти с символьного chunking на chunking по сообщениям
* **P1**: Учитывать time gap при сборке чанков
* **P1**: Маркировать в тексте `quote`, `forward`, автора сообщения и автора цитаты
* **P1**: Развести `page_content`, `dense_content`, `sparse_content`
* **P1**: Материализовать `sender_id` и `mentions` в индексируемый текст
* **P1**: Разбирать `file_snippets` и вытаскивать имя файла, mime и url
* **P1**: Разбирать `member_event` и превращать его в индексируемый текст
#### `index/requirements.txt` и/или `search/requirements.txt`
* **P4**: Добавить `razdel`
* **P4**: Добавить `pymorphy3`
* **P4**: Добавить `rapidfuzz`
---
### Человек 3 (Основной фокус на Docker и зависимостях):
#### `docker-compose.yml`
* **P3**: Привести локальный `docker-compose.yml` к схеме с `API_KEY`
* **P3**: Добавить `--platform linux/amd64` в сборку образов
#### `doc/` (новый файл с регрессионными вопросами)
* **P3**: Зафиксировать набор локальных тестовых вопросов для проверки регрессий
#### `search/requirements.txt`
* **P4**: Добавить `python-dateutil` или `dateparser`
* **P4**: Добавить `tenacity`
---
Таким образом, задачи равномерно распределены по 3 участникам с учётом сложности и области фокуса.

56
doc/upload_to_docker.md Normal file
View file

@ -0,0 +1,56 @@
Шаг 2. Настройка Docker
Registry для хранения образов будет доступен по адресу 83.166.249.64:5000. Поскольку он работает без TLS, необходимо добавить его в список insecure registries в настройках Docker.
Docker Desktop (macOS / Windows)
Откройте Docker Desktop -> Settings -> Docker Engine.
Добавьте в JSON-конфиг поле insecure-registries:
{
"insecure-registries": ["83.166.249.64:5000"]
}
Нажмите Apply & Restart.
CLI — Linux
Откройте файл с конфигурацией докер демона в режиме редактирования
sudo nano /etc/docker/daemon.json
Добавьте в JSON-конфиг поле insecure-registries:
{
"insecure-registries": ["83.166.249.64:5000"]
}
Перезапустите docker
sudo systemctl restart docker
Шаг 3. Логин в registry
Необходимо пройти аутентификацию в docker registry, используя логин и пароль, полученные на шаге 1.
docker login 83.166.249.64:5000 -u <login> -p <password>
Ожидаемый вывод: Login Succeeded.
Шаг 4. Сборка образов
Соберите образы своих Index Service и Search Service. Образы должны иметь тег, который имеет строгий формат:
Для Index Service - 83.166.249.64:5000/35230/index-service:latest
Для Search Service - 83.166.249.64:5000/35230/search-service:latest
Образы должны быть собраны под платформу linux/amd64
docker build --platform linux/amd64 -t 83.166.249.64:5000/35230/index-service:latest {path_to_index_service_dir}
docker build --platform linux/amd64 -t 83.166.249.64:5000/35230/search-service:latest {path_to_search_service_dir}
Шаг 5. Публикация образов в registry
docker push 83.166.249.64:5000/35230/index-service:latest
docker push 83.166.249.64:5000/35230/search-service:latest

View file

@ -2,39 +2,37 @@ services:
qdrant:
image: qdrant/qdrant:v1.14.1
ports:
- "6334:6333"
- "6333:6333"
qdrant-init:
image: curlimages/curl:8.12.1
env_file:
- .env
depends_on:
- qdrant
command:
- sh
- -c
- |
until curl -sf "$$QDRANT_URL/collections"; do
until curl -sf http://qdrant:6333/collections; do
sleep 1
done
if curl -sf "$$QDRANT_URL/collections/$$QDRANT_COLLECTION_NAME" >/dev/null; then
if curl -sf http://qdrant:6333/collections/evaluation >/dev/null; then
exit 0
fi
curl -sf -X PUT "$$QDRANT_URL/collections/$$QDRANT_COLLECTION_NAME" \
curl -sf -X PUT http://qdrant:6333/collections/evaluation \
-H 'Content-Type: application/json' \
-d "{
\"vectors\": {
\"$$QDRANT_DENSE_VECTOR_NAME\": {
\"size\": 1024,
\"distance\": \"Cosine\"
-d '{
"vectors": {
"dense": {
"size": 1024,
"distance": "Cosine"
}
},
\"sparse_vectors\": {
\"$$QDRANT_SPARSE_VECTOR_NAME\": {
\"modifier\": \"idf\"
"sparse_vectors": {
"sparse": {
"modifier": "idf"
}
}
}"
}'
restart: "no"
index:
@ -48,10 +46,18 @@ services:
search:
build:
context: ./search
env_file:
- .env
depends_on:
qdrant-init:
condition: service_completed_successfully
environment:
QDRANT_URL: http://qdrant:6333
QDRANT_COLLECTION_NAME: evaluation
QDRANT_DENSE_VECTOR_NAME: dense
QDRANT_SPARSE_VECTOR_NAME: sparse
EMBEDDINGS_DENSE_URL: ${EMBEDDINGS_DENSE_URL:-http://83.166.249.64:18001/embeddings}
RERANKER_URL: ${RERANKER_URL:-http://83.166.249.64:18001/score}
OPEN_API_LOGIN: ${OPEN_API_LOGIN:?set OPEN_API_LOGIN before docker compose up}
OPEN_API_PASSWORD: ${OPEN_API_PASSWORD:?set OPEN_API_PASSWORD before docker compose up}
ports:
- "8002:8000"

View file

@ -9,6 +9,11 @@ COPY main.py .
ENV HOST=0.0.0.0
ENV PORT=8000
ENV INDEX_CHUNK_MAX_CHARS=2200
ENV INDEX_MESSAGE_MAX_CHARS=1400
ENV INDEX_TEXT_SECTION_MAX_CHARS=1000
ENV INDEX_TIME_GAP_SECONDS=21600
ENV INDEX_CHUNK_OVERLAP_MESSAGES=1
ENV FASTEMBED_CACHE_PATH=/models/fastembed
ENV HF_HOME=/models/huggingface

View file

@ -1,26 +1,24 @@
LOGIN ?=
PASSWORD ?=
TEAM_ID ?= 35230
TEAM_ID ?=
DOCKER_REGISTRY_URL ?= 83.166.249.64:5000
PORT ?= 8000
IMAGE = $(DOCKER_REGISTRY_URL)/$(TEAM_ID)/index-service:latest
.PHONY: login build run push release
.PHONY: login build run push
login:
@: $(if $(LOGIN),,$(error LOGIN is required))
@: $(if $(PASSWORD),,$(error PASSWORD is required))
docker login $(DOCKER_REGISTRY_URL) -u $(LOGIN) -p $(PASSWORD)
@: $(if $(LOGIN),,$(error LOGIN is required for make login))
@: $(if $(PASSWORD),,$(error PASSWORD is required for make login))
docker login $(DOCKER_REGISTRY_URL) -u $(LOGIN) -p $(PASSWORD)
build:
docker build --platform linux/amd64 -t $(IMAGE) ./
@: $(if $(TEAM_ID),,$(error TEAM_ID is required for make build))
docker build -t $(IMAGE) ./
run: build
docker run --rm -p $(PORT):8000 $(IMAGE)
push:
push: login build
docker push $(IMAGE)
release: build push
@echo "index-service pushed → $(IMAGE)"

View file

@ -1,107 +0,0 @@
"""Message-based chunking with window by count, length, and time gap."""
from cleaning import CleanedMessage, clean_message
from rendering import render_dense_content, render_page_content, render_sparse_content
from index_schemas import IndexAPIItem, Message
WINDOW_MAX_MESSAGES = 5
WINDOW_MAX_CHARS = 512
TIME_GAP_SECONDS = 3600
OVERLAP_MESSAGES = 2
def _clean_all(messages: list[Message]) -> list[CleanedMessage]:
cleaned = [clean_message(m) for m in messages if not m.is_system and not m.is_hidden]
return [c for c in cleaned if not c.is_empty]
def _render_chunk(
overlap: list[CleanedMessage],
window: list[CleanedMessage],
) -> IndexAPIItem:
page_lines: list[str] = []
dense_lines: list[str] = []
sparse_tokens: list[str] = []
for msg in overlap + window:
page = render_page_content(msg)
dense = render_dense_content(msg)
sparse = render_sparse_content(msg)
if page:
page_lines.append(page)
if dense:
dense_lines.append(dense)
if sparse:
sparse_tokens.append(sparse)
return IndexAPIItem(
page_content="\n".join(page_lines),
dense_content="\n".join(dense_lines),
sparse_content=" ".join(sparse_tokens),
message_ids=[msg.id for msg in window],
)
def _split_windows(messages: list[CleanedMessage]) -> list[list[CleanedMessage]]:
"""Split cleaned messages into windows respecting count, length, and time gap."""
if not messages:
return []
windows: list[list[CleanedMessage]] = []
current: list[CleanedMessage] = []
current_chars = 0
for msg in messages:
msg_text = render_page_content(msg)
msg_chars = len(msg_text)
time_break = (
current
and (msg.time - current[-1].time) > TIME_GAP_SECONDS
)
size_break = (
current
and (
len(current) >= WINDOW_MAX_MESSAGES
or current_chars + msg_chars > WINDOW_MAX_CHARS
)
)
if time_break or size_break:
if current:
windows.append(current)
current = [msg]
current_chars = msg_chars
else:
current.append(msg)
current_chars += msg_chars
if current:
windows.append(current)
return windows
def build_chunks(
overlap_messages: list[Message],
new_messages: list[Message],
) -> list[IndexAPIItem]:
clean_overlap = _clean_all(overlap_messages)
clean_new = _clean_all(new_messages)
if not clean_new:
return []
overlap_tail = clean_overlap[-OVERLAP_MESSAGES:] if clean_overlap else []
windows = _split_windows(clean_new)
result: list[IndexAPIItem] = []
prev_window_tail: list[CleanedMessage] = overlap_tail
for window in windows:
chunk = _render_chunk(prev_window_tail, window)
if chunk.message_ids:
result.append(chunk)
prev_window_tail = window[-OVERLAP_MESSAGES:]
return result

View file

@ -1,157 +0,0 @@
"""Local message cleaning and normalization. No external API calls."""
import json
import re
from typing import Any
_ZERO_WIDTH = re.compile(r"[\u200b\u200c\u200d\ufeff]")
_MULTI_NEWLINE = re.compile(r"\n{3,}")
_MULTI_SPACE = re.compile(r"[ \t]{2,}")
def normalize_unicode(text: str) -> str:
text = _ZERO_WIDTH.sub("", text)
text = text.replace("\r\n", "\n").replace("\r", "\n")
text = _MULTI_SPACE.sub(" ", text)
text = _MULTI_NEWLINE.sub("\n\n", text)
return text.strip()
def _safe_normalize(text: str | None) -> str:
if not text:
return ""
return normalize_unicode(str(text))
def parse_file_snippets(raw: str) -> list[dict[str, Any]]:
if not raw or not raw.strip():
return []
try:
parsed = json.loads(raw)
if isinstance(parsed, list):
return parsed
if isinstance(parsed, dict):
return [parsed]
return []
except (json.JSONDecodeError, ValueError):
return []
def extract_file_info(snippet: dict[str, Any]) -> dict[str, str]:
return {
"name": str(snippet.get("name") or ""),
"mime": str(snippet.get("mime") or ""),
"url": str(snippet.get("original_url") or ""),
"date": str(snippet.get("date_create") or ""),
}
def normalize_member_event(event: dict[str, Any] | None) -> str:
if not event:
return ""
event_type = str(event.get("type") or "unknown_event")
members = event.get("members") or []
if event_type == "addMembers" and members:
joined = ", ".join(str(m) for m in members)
return f"[system: {joined} added to chat]"
payload_str = json.dumps(event, ensure_ascii=False)
return f"[system: {event_type} {payload_str}]"
def normalize_part(part: dict[str, Any]) -> dict[str, str]:
"""Normalize a single message part by its mediaType."""
media_type = str(part.get("mediaType") or "text")
text = _safe_normalize(part.get("text"))
if media_type == "text":
return {"type": "text", "text": text}
if media_type == "quote":
sender = _safe_normalize(part.get("sn") or part.get("sender_id") or "")
label = f"[quote from {sender}]" if sender else "[quote]"
return {"type": "quote", "text": f"{label}: {text}" if text else label}
if media_type == "forward":
origin = _safe_normalize(part.get("sn") or "")
label = f"[forwarded from {origin}]" if origin else "[forwarded]"
return {"type": "forward", "text": f"{label}: {text}" if text else label}
label = f"[{media_type}]"
return {"type": media_type, "text": f"{label}: {text}" if text else label}
class CleanedMessage:
__slots__ = (
"id",
"sender_id",
"time",
"thread_sn",
"text",
"parts",
"mentions",
"member_event_text",
"file_info",
"is_system",
"is_forward",
"is_quote",
"is_empty",
)
def __init__(
self,
id: str,
sender_id: str,
time: int,
thread_sn: str | None,
text: str,
parts: list[dict[str, str]],
mentions: list[str],
member_event_text: str,
file_info: list[dict[str, str]],
is_system: bool,
is_forward: bool,
is_quote: bool,
):
self.id = id
self.sender_id = sender_id
self.time = time
self.thread_sn = thread_sn
self.text = text
self.parts = parts
self.mentions = mentions
self.member_event_text = member_event_text
self.file_info = file_info
self.is_system = is_system
self.is_forward = is_forward
self.is_quote = is_quote
self.is_empty = not (text or parts or member_event_text or file_info)
def clean_message(msg: Any) -> CleanedMessage:
"""Extract and normalize all signals from a raw Message object."""
text = _safe_normalize(msg.text)
parts = [normalize_part(p) for p in (msg.parts or []) if isinstance(p, dict)]
parts = [p for p in parts if p.get("text")]
mentions = [_safe_normalize(m) for m in (msg.mentions or []) if m]
member_event_text = normalize_member_event(msg.member_event)
file_info: list[dict[str, str]] = []
for snippet in parse_file_snippets(msg.file_snippets):
info = extract_file_info(snippet)
if any(info.values()):
file_info.append(info)
return CleanedMessage(
id=msg.id,
sender_id=_safe_normalize(msg.sender_id),
time=msg.time,
thread_sn=msg.thread_sn,
text=text,
parts=parts,
mentions=mentions,
member_event_text=member_event_text,
file_info=file_info,
is_system=bool(msg.is_system),
is_forward=bool(msg.is_forward),
is_quote=bool(msg.is_quote),
)

View file

@ -1,59 +0,0 @@
from typing import Any
from pydantic import BaseModel
class Chat(BaseModel):
id: str
name: str
sn: str
type: str
is_public: bool | None = None
members_count: int | None = None
members: list[dict[str, Any]] | None = None
class Message(BaseModel):
id: str
thread_sn: str | None = None
time: int
text: str
sender_id: str
file_snippets: str
parts: list[dict[str, Any]] | None = None
mentions: list[str] | None = None
member_event: dict[str, Any] | None = None
is_system: bool
is_hidden: bool
is_forward: bool
is_quote: bool
class ChatData(BaseModel):
chat: Chat
overlap_messages: list[Message]
new_messages: list[Message]
class IndexAPIRequest(BaseModel):
data: ChatData
class IndexAPIItem(BaseModel):
page_content: str
dense_content: str
sparse_content: str
message_ids: list[str]
class IndexAPIResponse(BaseModel):
results: list[IndexAPIItem]
class SparseEmbeddingRequest(BaseModel):
texts: list[str]
class SparseVector(BaseModel):
indices: list[int]
values: list[float]

View file

@ -1,27 +1,32 @@
import asyncio
import logging
import os
import json
import re
from dataclasses import dataclass
from datetime import datetime, timezone
from functools import lru_cache
from typing import Any
import asyncio
from fastapi import FastAPI, Request
from fastapi.exceptions import RequestValidationError
from fastapi.responses import JSONResponse
from pydantic import BaseModel
# Ваш сервис должен считывать эти переменные из окружения (env), так как проверяющая система управляет ими
HOST = os.getenv("HOST", "0.0.0.0")
PORT = int(os.getenv("PORT", "8000"))
UVICORN_WORKERS = 8
PORT = int(os.getenv("PORT", "8004"))
logging.basicConfig(level=os.getenv("LOG_LEVEL", "INFO"))
logger = logging.getLogger("index-service")
# Модель данных, которую мы предоставляем и рассчитываем получать от вас
class Chat(BaseModel):
id: str
name: str
sn: str
type: str
type: str # group, channel, private
is_public: bool | None = None
members_count: int | None = None
members: list[dict[str, Any]] | None = None
@ -53,6 +58,10 @@ class IndexAPIRequest(BaseModel):
data: ChatData
# dense_content будет передан в dense embedding модель для построения семантического вектора.
# sparse_content будет передан в sparse модель для построения разреженного индекса "по словам".
# Можно оставить dense_content и sparse_content равными page_content,
# а можно формировать для них разные версии текста.
class IndexAPIItem(BaseModel):
page_content: str
dense_content: str
@ -73,122 +82,534 @@ class SparseVector(BaseModel):
values: list[float]
CHUNK_SIZE = 256
OVERLAP_SIZE = 128
SPARSE_MODEL_NAME = "Qdrant/bm25"
FASTEMBED_CACHE_PATH = "/models/fastembed"
def render_message(message: Message) -> str:
parts_list: list[str] = []
if message.sender_id:
sender_name = message.sender_id.split("@")[0].replace(".", " ")
parts_list.append(f"[{sender_name}]:")
if message.text:
parts_list.append(message.text)
if message.parts:
for part in message.parts:
media_type = part.get("mediaType", "text")
part_text = part.get("text")
if isinstance(part_text, str) and part_text:
if media_type == "forward":
parts_list.append(f"[Пересланное]: {part_text}")
elif media_type == "quote":
parts_list.append(f"[Цитата]: {part_text}")
else:
parts_list.append(part_text)
if message.file_snippets:
parts_list.append(f"[Файл]: {message.file_snippets}")
return " ".join(parts_list).strip()
def build_chunks(
chat: Chat,
overlap_messages: list[Message],
new_messages: list[Message],
) -> list[IndexAPIItem]:
new_messages = [m for m in new_messages if not m.is_system and not m.is_hidden]
overlap_messages = [m for m in overlap_messages if not m.is_system and not m.is_hidden]
result: list[IndexAPIItem] = []
def build_text_and_ranges(messages: list[Message]) -> tuple[str, list[tuple[int, int, str]]]:
text_parts: list[str] = []
message_ranges: list[tuple[int, int, str]] = []
position = 0
for index, message in enumerate(messages):
text = render_message(message)
if not text:
continue
if index > 0 and text_parts:
text_parts.append("\n")
position += 1
start = position
text_parts.append(text)
position += len(text)
message_ranges.append((start, position, message.id))
return "".join(text_parts), message_ranges
def slice_tail(text: str, tail_size: int) -> str:
if tail_size <= 0:
return ""
tail_start = max(0, len(text) - tail_size)
return text[tail_start:]
overlap_text, _ = build_text_and_ranges(overlap_messages)
previous_chunk_text = slice_tail(overlap_text, OVERLAP_SIZE)
new_text, new_message_ranges = build_text_and_ranges(new_messages)
for start in range(0, len(new_text), CHUNK_SIZE):
chunk_body = new_text[start: start + CHUNK_SIZE]
if not chunk_body:
continue
chunk_body_ranges = [
(
max(message_start, start) - start,
min(message_end, start + len(chunk_body)) - start,
message_id,
)
for message_start, message_end, message_id in new_message_ranges
if message_end > start and message_start < start + len(chunk_body)
]
chunk_overlap = previous_chunk_text
chunk_text = chunk_overlap
if chunk_text and chunk_body:
chunk_text += "\n"
chunk_text += chunk_body
dense_text = f"[{chat.name}] {chunk_text}"
sparse_text = chunk_body
result.append(
IndexAPIItem(
page_content=chunk_text,
dense_content=dense_text,
sparse_content=sparse_text,
message_ids=[message_id for _, _, message_id in chunk_body_ranges],
)
)
previous_chunk_text = slice_tail(chunk_text, OVERLAP_SIZE)
return result
class SparseEmbeddingResponse(BaseModel):
vectors: list[SparseVector]
app = FastAPI(title="Index Service", version="0.1.0")
# Ваша внутренняя логика построения чанков. Можете делать всё, что посчитаете нужным.
# Текущий код минимальный пример
CHUNK_MAX_CHARS = int(os.getenv("INDEX_CHUNK_MAX_CHARS", "2200"))
MESSAGE_MAX_CHARS = int(os.getenv("INDEX_MESSAGE_MAX_CHARS", "1400"))
TEXT_SECTION_MAX_CHARS = int(os.getenv("INDEX_TEXT_SECTION_MAX_CHARS", "1000"))
TIME_GAP_SECONDS = int(os.getenv("INDEX_TIME_GAP_SECONDS", str(6 * 60 * 60)))
CHUNK_OVERLAP_MESSAGES = int(os.getenv("INDEX_CHUNK_OVERLAP_MESSAGES", "1"))
SPARSE_MODEL_NAME = "Qdrant/bm25"
FASTEMBED_CACHE_PATH = "/models/fastembed"
# Важная переманная, которая позволяет вычислять sparse вектор в несколько ядер. Не рекомендуется изменять.
UVICORN_WORKERS = 8
@dataclass(frozen=True)
class RenderedText:
page: str
dense: str
sparse: str
@dataclass(frozen=True)
class RenderedUnit:
message_id: str
time: int
text: RenderedText
@dataclass(frozen=True)
class ChunkUnit:
unit: RenderedUnit
is_new: bool
def clean_text(value: Any) -> str:
if not isinstance(value, str):
return ""
text = (
value.replace("\r\n", "\n")
.replace("\r", "\n")
.replace("\u200b", " ")
.replace("\xa0", " ")
)
lines = [re.sub(r"[ \t]+", " ", line).strip() for line in text.split("\n")]
result: list[str] = []
previous_blank = False
for line in lines:
if not line:
if result and not previous_blank:
result.append("")
previous_blank = True
continue
result.append(line)
previous_blank = False
return "\n".join(result).strip()
def unique_preserve_order(values: list[str]) -> list[str]:
seen: set[str] = set()
result: list[str] = []
for value in values:
if value and value not in seen:
seen.add(value)
result.append(value)
return result
def format_time(timestamp: int) -> str:
return datetime.fromtimestamp(timestamp, tz=timezone.utc).isoformat().replace("+00:00", "Z")
def format_optional_time(value: Any) -> str:
if value is None or value == "":
return ""
try:
return format_time(int(value))
except (TypeError, ValueError, OSError, OverflowError):
return clean_text(str(value))
def lexical_terms(value: str) -> str:
return clean_text(re.sub(r"[^0-9A-Za-zА-Яа-яЁё]+", " ", value))
def split_long_text(text: str, max_chars: int) -> list[str]:
text = clean_text(text)
if not text:
return []
if len(text) <= max_chars:
return [text]
paragraphs = [item.strip() for item in re.split(r"\n{2,}", text) if item.strip()]
pieces: list[str] = []
current = ""
def append_current() -> None:
nonlocal current
if current:
pieces.append(current)
current = ""
def split_oversized(paragraph: str) -> list[str]:
words = paragraph.split()
result: list[str] = []
part = ""
for word in words:
if not part:
part = word
continue
if len(part) + 1 + len(word) <= max_chars:
part += " " + word
else:
result.append(part)
part = word
if part:
result.append(part)
return result
for paragraph in paragraphs:
candidates = [paragraph] if len(paragraph) <= max_chars else split_oversized(paragraph)
for candidate in candidates:
separator = "\n\n" if current else ""
if current and len(current) + len(separator) + len(candidate) > max_chars:
append_current()
current = candidate if not current else current + separator + candidate
append_current()
return pieces
def combine_texts(items: list[RenderedText], separator: str = "\n") -> RenderedText:
return RenderedText(
page=separator.join(item.page for item in items if item.page).strip(),
dense=separator.join(item.dense for item in items if item.dense).strip(),
sparse=separator.join(item.sparse for item in items if item.sparse).strip(),
)
def rendered_length(text: RenderedText) -> int:
return max(len(text.page), len(text.dense), len(text.sparse))
def message_header(message: Message) -> RenderedText:
timestamp = format_time(message.time)
mentions = unique_preserve_order(message.mentions or [])
page_lines = [
f"author: {message.sender_id}",
f"time: {timestamp}",
]
dense_lines = [
f"author: {message.sender_id}",
f"message_time: {timestamp}",
]
sparse_terms = [
"author",
message.sender_id,
lexical_terms(message.sender_id),
timestamp,
]
if message.thread_sn:
page_lines.append(f"thread: {message.thread_sn}")
dense_lines.append(f"thread: {message.thread_sn}")
sparse_terms.extend(["thread", message.thread_sn, lexical_terms(message.thread_sn)])
if mentions:
mentions_text = ", ".join(mentions)
page_lines.append(f"mentions: {mentions_text}")
dense_lines.append(f"mentions: {mentions_text}")
sparse_terms.extend(["mentions", *mentions, *(lexical_terms(item) for item in mentions)])
flags = []
if message.is_system:
flags.append("system")
if message.is_forward:
flags.append("forward")
if message.is_quote:
flags.append("quote")
if message.is_hidden:
flags.append("hidden")
if flags:
flags_text = ", ".join(flags)
dense_lines.append(f"message_flags: {flags_text}")
sparse_terms.extend(flags)
return RenderedText(
page="\n".join(page_lines),
dense="\n".join(dense_lines),
sparse=" ".join(term for term in sparse_terms if term),
)
def render_plain_sections(text: str, label: str = "text") -> list[RenderedText]:
sections: list[RenderedText] = []
pieces = split_long_text(text, TEXT_SECTION_MAX_CHARS)
for index, piece in enumerate(pieces):
suffix = f" part {index + 1}/{len(pieces)}" if len(pieces) > 1 else ""
sections.append(
RenderedText(
page=piece,
dense=f"{label}{suffix}: {piece}",
sparse=f"{label} {piece}",
)
)
return sections
def render_part_sections(part: dict[str, Any]) -> list[RenderedText]:
text = clean_text(part.get("text"))
if not text:
return []
media_type = clean_text(part.get("mediaType") or part.get("type") or "text").lower()
source = clean_text(part.get("sn"))
part_time = part.get("time")
source_bits = []
if source:
source_bits.append(f"source: {source}")
formatted_part_time = format_optional_time(part_time)
if formatted_part_time:
source_bits.append(f"source_time: {formatted_part_time}")
source_text = ", ".join(source_bits)
if media_type == "quote":
label = f"quote from {source}" if source else "quote"
dense_label = f"quote; {source_text}" if source_text else "quote"
sparse_prefix = f"quote цитата {source} {lexical_terms(source)}"
elif media_type == "forward":
label = f"forwarded from {source}" if source else "forwarded"
dense_label = f"forwarded_message; {source_text}" if source_text else "forwarded_message"
sparse_prefix = f"forward forwarded_message пересланное {source} {lexical_terms(source)}"
else:
label = "text"
dense_label = "text"
sparse_prefix = "text"
sections: list[RenderedText] = []
pieces = split_long_text(text, TEXT_SECTION_MAX_CHARS)
for index, piece in enumerate(pieces):
suffix = f" part {index + 1}/{len(pieces)}" if len(pieces) > 1 else ""
page_prefix = f"{label}{suffix}:"
sections.append(
RenderedText(
page=f"{page_prefix}\n{piece}" if media_type in {"quote", "forward"} else piece,
dense=f"{dense_label}{suffix}: {piece}",
sparse=f"{sparse_prefix} {piece}",
)
)
return sections
def render_file_sections(raw_snippets: str) -> list[RenderedText]:
raw_snippets = clean_text(raw_snippets)
if not raw_snippets:
return []
try:
parsed = json.loads(raw_snippets)
except json.JSONDecodeError:
return [
RenderedText(
page=f"file_snippet: {raw_snippets}",
dense=f"file_snippet: {raw_snippets}",
sparse=f"file file_snippet {raw_snippets}",
)
]
snippets = parsed if isinstance(parsed, list) else [parsed]
sections: list[RenderedText] = []
for snippet in snippets:
if not isinstance(snippet, dict):
text = clean_text(str(snippet))
sections.append(RenderedText(page=f"file: {text}", dense=f"file: {text}", sparse=f"file {text}"))
continue
name = clean_text(snippet.get("name"))
mime = clean_text(snippet.get("mime"))
url = clean_text(snippet.get("original_url") or snippet.get("url"))
owner = clean_text(snippet.get("uid"))
created = clean_text(snippet.get("date_create"))
file_bits = [
f"name: {name}" if name else "",
f"mime: {mime}" if mime else "",
f"url: {url}" if url else "",
f"owner: {owner}" if owner else "",
f"created: {created}" if created else "",
]
file_text = ", ".join(bit for bit in file_bits if bit)
sparse_terms = " ".join(
term
for term in [
"file",
"attachment",
"document",
name,
lexical_terms(name),
mime,
url,
owner,
lexical_terms(owner),
created,
]
if term
)
sections.append(
RenderedText(
page=f"file: {file_text}",
dense=f"file: {file_text}",
sparse=sparse_terms,
)
)
return sections
def render_member_event(message: Message) -> list[RenderedText]:
event = message.member_event
if not event:
return []
event_type = clean_text(event.get("type") or "member_event")
members_raw = event.get("members")
members = [clean_text(item) for item in members_raw] if isinstance(members_raw, list) else []
members = unique_preserve_order([item for item in members if item])
if members:
members_text = ", ".join(members)
else:
members_text = " ".join(clean_text(str(value)) for value in event.values() if value)
page = f"system_event: {event_type}; actor: {message.sender_id}; members: {members_text}"
dense = (
f"system_event: {event_type}; action: add or update chat members; "
f"actor: {message.sender_id}; members: {members_text}"
)
sparse = " ".join(
term
for term in [
"system_event",
"member_event",
event_type,
"addMembers",
"добавление участников",
message.sender_id,
lexical_terms(message.sender_id),
members_text,
lexical_terms(members_text),
]
if term
)
return [RenderedText(page=page, dense=dense, sparse=sparse)]
def render_message_sections(message: Message) -> list[RenderedText]:
sections: list[RenderedText] = []
sections.extend(render_plain_sections(message.text, "message_text"))
for part in message.parts or []:
if isinstance(part, dict):
sections.extend(render_part_sections(part))
sections.extend(render_file_sections(message.file_snippets))
sections.extend(render_member_event(message))
return sections
def render_message_units(message: Message) -> list[RenderedUnit]:
sections = render_message_sections(message)
if not sections:
return []
header = message_header(message)
units: list[RenderedUnit] = []
current: list[RenderedText] = []
def flush() -> None:
nonlocal current
if not current:
return
text = combine_texts([header, *current])
units.append(RenderedUnit(message_id=message.id, time=message.time, text=text))
current = []
for section in sections:
candidate = combine_texts([header, *current, section])
if current and rendered_length(candidate) > MESSAGE_MAX_CHARS:
flush()
current.append(section)
flush()
return units
def chunk_text(items: list[ChunkUnit]) -> RenderedText:
return RenderedText(
page="\n\n".join(item.unit.text.page for item in items if item.unit.text.page).strip(),
dense="\n\n".join(item.unit.text.dense for item in items if item.unit.text.dense).strip(),
sparse="\n\n".join(item.unit.text.sparse for item in items if item.unit.text.sparse).strip(),
)
def chunk_length(items: list[ChunkUnit]) -> int:
return rendered_length(chunk_text(items))
def trim_context(context: list[RenderedUnit], unit: RenderedUnit) -> list[RenderedUnit]:
result = context[-CHUNK_OVERLAP_MESSAGES:] if CHUNK_OVERLAP_MESSAGES > 0 else []
items = [ChunkUnit(item, False) for item in result] + [ChunkUnit(unit, True)]
while result and chunk_length(items) > CHUNK_MAX_CHARS:
result = result[1:]
items = [ChunkUnit(item, False) for item in result] + [ChunkUnit(unit, True)]
return result
def build_chunks(
overlap_messages: list[Message],
new_messages: list[Message],
) -> list[IndexAPIItem]:
result: list[IndexAPIItem] = []
overlap_units = [
unit
for message in overlap_messages
for unit in render_message_units(message)
]
new_units = [
unit
for message in new_messages
for unit in render_message_units(message)
]
current: list[ChunkUnit] = []
last_new_time: int | None = None
def flush_current() -> None:
nonlocal current
if not current:
return
message_ids = unique_preserve_order(
[item.unit.message_id for item in current if item.is_new]
)
if not message_ids:
current = []
return
text = chunk_text(current)
result.append(
IndexAPIItem(
page_content=text.page,
dense_content=text.dense,
sparse_content=text.sparse,
message_ids=message_ids,
)
)
current = []
def request_overlap_context(unit: RenderedUnit) -> list[RenderedUnit]:
close_units = [
item
for item in overlap_units
if abs(unit.time - item.time) <= TIME_GAP_SECONDS
]
return trim_context(close_units, unit)
for unit in new_units:
if not current:
context = request_overlap_context(unit)
current = [ChunkUnit(item, False) for item in context]
current.append(ChunkUnit(unit, True))
last_new_time = unit.time
continue
gap = abs(unit.time - last_new_time) if last_new_time is not None else 0
candidate = [*current, ChunkUnit(unit, True)]
should_split = gap > TIME_GAP_SECONDS or chunk_length(candidate) > CHUNK_MAX_CHARS
if should_split:
previous_new_units = [item.unit for item in current if item.is_new]
context = (
trim_context(previous_new_units, unit)
if gap <= TIME_GAP_SECONDS
else request_overlap_context(unit)
)
flush_current()
current = [ChunkUnit(item, False) for item in context]
current.append(ChunkUnit(unit, True))
else:
current.append(ChunkUnit(unit, True))
last_new_time = unit.time
flush_current()
return result
# Ваш сервис должен имплементировать оба этих метода
@app.get("/health")
async def health() -> dict[str, str]:
return {"status": "ok"}
@ -198,7 +619,6 @@ async def health() -> dict[str, str]:
async def index(payload: IndexAPIRequest) -> IndexAPIResponse:
return IndexAPIResponse(
results=build_chunks(
payload.data.chat,
payload.data.overlap_messages,
payload.data.new_messages,
)
@ -209,13 +629,20 @@ async def index(payload: IndexAPIRequest) -> IndexAPIResponse:
def get_sparse_model():
from fastembed import SparseTextEmbedding
logger.info("Loading sparse model %s from cache %s", SPARSE_MODEL_NAME, FASTEMBED_CACHE_PATH)
# можете делать любой вектор, который будет совместим с вашим поиском в Qdrant
# помните об ограничении времени выполнения вашей работы в тестирующей системе
logger.info(
"Loading sparse model %s from cache %s",
SPARSE_MODEL_NAME,
FASTEMBED_CACHE_PATH,
)
return SparseTextEmbedding(model_name=SPARSE_MODEL_NAME)
def embed_sparse_texts(texts: list[str]) -> list[dict]:
def embed_sparse_texts(texts: list[str]) -> list[SparseVector]:
model = get_sparse_model()
vectors = []
vectors: list[dict[str, list[int] | list[float]]] = []
for item in model.embed(texts):
vectors.append(
{
@ -223,27 +650,37 @@ def embed_sparse_texts(texts: list[str]) -> list[dict]:
"values": item.values.tolist(),
}
)
return vectors
@app.post("/sparse_embedding")
async def sparse_embedding(payload: SparseEmbeddingRequest) -> dict[str, Any]:
# Проверяющая система вызывает этот endpoint при создании коллекции
vectors = await asyncio.to_thread(embed_sparse_texts, payload.texts)
return {"vectors": vectors}
# красивая обработка ошибок
@app.exception_handler(Exception)
async def exception_handler(request: Request, exc: Exception) -> JSONResponse:
logger.exception(exc)
if isinstance(exc, RequestValidationError):
return JSONResponse(status_code=422, content={"detail": exc.errors()})
return JSONResponse(status_code=500, content={"detail": str(exc)})
def main() -> None:
import uvicorn
uvicorn.run("main:app", host=HOST, port=PORT, reload=False, workers=UVICORN_WORKERS)
uvicorn.run(
"main:app",
host=HOST,
port=PORT,
reload=False,
workers=UVICORN_WORKERS,
)
if __name__ == "__main__":

View file

@ -1,102 +0,0 @@
"""Render cleaned messages into page_content, dense_content, sparse_content."""
import datetime
from cleaning import CleanedMessage
def _format_time(ts: int) -> str:
try:
return datetime.datetime.fromtimestamp(ts, tz=datetime.timezone.utc).strftime("%Y-%m-%d %H:%M")
except (OSError, OverflowError, ValueError):
return str(ts)
def render_page_content(msg: CleanedMessage) -> str:
"""Human-readable text for the chunk payload."""
lines: list[str] = []
prefix = f"{msg.sender_id}:"
if msg.text:
lines.append(f"{prefix} {msg.text}")
elif not msg.parts and not msg.member_event_text:
lines.append(prefix)
for part in msg.parts:
lines.append(part["text"])
if msg.member_event_text:
lines.append(msg.member_event_text)
for fi in msg.file_info:
name = fi.get("name") or fi.get("url") or "file"
lines.append(f"[attachment: {name}]")
return "\n".join(lines)
def render_dense_content(msg: CleanedMessage) -> str:
"""Text optimized for semantic dense embedding: role markers + full context."""
lines: list[str] = []
ts = _format_time(msg.time)
header_parts = [f"[{ts}]", f"sender:{msg.sender_id}"]
if msg.is_forward:
header_parts.append("type:forward")
if msg.is_quote:
header_parts.append("type:quote")
if msg.is_system:
header_parts.append("type:system")
if msg.thread_sn:
header_parts.append(f"thread:{msg.thread_sn}")
lines.append(" ".join(header_parts))
if msg.text:
lines.append(msg.text)
for part in msg.parts:
lines.append(part["text"])
if msg.member_event_text:
lines.append(msg.member_event_text)
for fi in msg.file_info:
fi_parts = []
if fi.get("name"):
fi_parts.append(fi["name"])
if fi.get("mime"):
fi_parts.append(fi["mime"])
if fi.get("url"):
fi_parts.append(fi["url"])
if fi_parts:
lines.append(f"[file: {' '.join(fi_parts)}]")
if msg.mentions:
lines.append("mentions: " + ", ".join(msg.mentions))
return "\n".join(lines)
def render_sparse_content(msg: CleanedMessage) -> str:
"""Keyword-heavy text for sparse/BM25 embedding."""
tokens: list[str] = []
tokens.append(msg.sender_id)
tokens.extend(msg.mentions)
if msg.text:
tokens.append(msg.text)
for part in msg.parts:
tokens.append(part["text"])
if msg.member_event_text:
tokens.append(msg.member_event_text)
for fi in msg.file_info:
for key in ("name", "mime", "url"):
val = fi.get(key)
if val:
tokens.append(val)
return " ".join(tokens)

View file

@ -1,31 +0,0 @@
import logging
import os
from functools import lru_cache
from index_schemas import SparseVector
SPARSE_MODEL_NAME = "Qdrant/bm25"
FASTEMBED_CACHE_PATH = "/models/fastembed"
logger = logging.getLogger("index-service")
@lru_cache(maxsize=1)
def get_sparse_model():
from fastembed import SparseTextEmbedding
logger.info("Loading sparse model %s from cache %s", SPARSE_MODEL_NAME, FASTEMBED_CACHE_PATH)
return SparseTextEmbedding(model_name=SPARSE_MODEL_NAME)
def embed_sparse_texts(texts: list[str]) -> list[SparseVector]:
model = get_sparse_model()
result: list[SparseVector] = []
for item in model.embed(texts):
result.append(
SparseVector(
indices=[int(i) for i in item.indices.tolist()],
values=[float(v) for v in item.values.tolist()],
)
)
return result

View file

@ -1,41 +0,0 @@
---
## 3) `bin/codex-start` чтобы всё поднималось само
`SKILL.md` сам по себе **не умеет магически запускать Codex при открытии терминала**. Для этого нужен обычный wrapper-скрипт. Вот рабочий вариант:
```bash
#!/usr/bin/env bash
set -euo pipefail
ROOT="$(git rev-parse --show-toplevel 2>/dev/null || pwd)"
cd "$ROOT"
mkdir -p .ai_update/sessions
touch .ai_update/current_status.md
touch .ai_update/handoff.md
touch .ai_update/changelog.md
touch .ai_update/touched_files.md
COMPOSE_FILE=""
if [ -f docker-compose.yml ]; then
COMPOSE_FILE="docker-compose.yml"
elif [ -f compose.yaml ]; then
COMPOSE_FILE="compose.yaml"
elif [ -f compose.yml ]; then
COMPOSE_FILE="compose.yml"
fi
if [ -n "$COMPOSE_FILE" ] && command -v docker >/dev/null 2>&1; then
if ! docker compose ps --status running >/dev/null 2>&1; then
docker compose up -d --build || true
else
RUNNING_COUNT="$(docker compose ps --status running --services 2>/dev/null | wc -l | tr -d ' ')"
if [ "${RUNNING_COUNT:-0}" = "0" ]; then
docker compose up -d --build || true
fi
fi
fi
exec codex "$@"

View file

@ -1,6 +1,6 @@
LOGIN ?=
PASSWORD ?=
TEAM_ID ?= 35230
TEAM_ID ?=
DOCKER_REGISTRY_URL ?= 83.166.249.64:5000
PORT ?= 8000
QDRANT_URL ?=
@ -16,15 +16,16 @@ REQUIRED_RUN_VARS := QDRANT_URL EMBEDDINGS_DENSE_URL API_KEY RERANKER_URL
IMAGE = $(DOCKER_REGISTRY_URL)/$(TEAM_ID)/search-service:latest
.PHONY: login build run push release
.PHONY: login build run push check-run-env
login:
@: $(if $(LOGIN),,$(error LOGIN is required))
@: $(if $(PASSWORD),,$(error PASSWORD is required))
@: $(if $(LOGIN),,$(error LOGIN is required for make login))
@: $(if $(PASSWORD),,$(error PASSWORD is required for make login))
docker login $(DOCKER_REGISTRY_URL) -u $(LOGIN) -p $(PASSWORD)
build:
docker build --platform linux/amd64 -t $(IMAGE) ./
@: $(if $(TEAM_ID),,$(error TEAM_ID is required for make build))
docker build -t $(IMAGE) ./
run: build
@: $(foreach var,$(REQUIRED_RUN_VARS),$(if $($(var)),,$(error $(var) is required for make run)))
@ -40,8 +41,5 @@ run: build
-e QDRANT_SPARSE_VECTOR_NAME=$(QDRANT_SPARSE_VECTOR_NAME) \
$(IMAGE)
push:
push: login build
docker push $(IMAGE)
release: build push
@echo "search-service pushed → $(IMAGE)"

View file

View file

@ -1,23 +0,0 @@
from typing import Any
from config import TOP_K
from retrieval import extract_message_ids
def aggregate_message_ids(
reranked_head: list[Any],
retrieval_tail: list[Any],
) -> list[str]:
"""Deduplicate and collect top-K message_ids, reranked head first."""
seen: set[str] = set()
result: list[str] = []
for point in reranked_head + retrieval_tail:
for mid in extract_message_ids(point):
if mid not in seen:
seen.add(mid)
result.append(mid)
if len(result) >= TOP_K:
return result
return result

View file

@ -1,56 +0,0 @@
import logging
import os
from typing import Any
EMBEDDINGS_DENSE_MODEL = "Qwen/Qwen3-Embedding-0.6B"
SPARSE_MODEL_NAME = "Qdrant/bm25"
RERANKER_MODEL = "nvidia/llama-nemotron-rerank-1b-v2"
HOST = os.getenv("HOST", "0.0.0.0")
PORT = int(os.getenv("PORT", "8003"))
API_KEY = os.getenv("API_KEY")
EMBEDDINGS_DENSE_URL = os.getenv("EMBEDDINGS_DENSE_URL")
RERANKER_URL = os.getenv("RERANKER_URL")
QDRANT_URL = os.getenv("QDRANT_URL")
QDRANT_COLLECTION_NAME = os.getenv("QDRANT_COLLECTION_NAME", "evaluation")
QDRANT_DENSE_VECTOR_NAME = os.getenv("QDRANT_DENSE_VECTOR_NAME", "dense")
QDRANT_SPARSE_VECTOR_NAME = os.getenv("QDRANT_SPARSE_VECTOR_NAME", "sparse")
OPEN_API_LOGIN = os.getenv("OPEN_API_LOGIN")
OPEN_API_PASSWORD = os.getenv("OPEN_API_PASSWORD")
DENSE_PREFETCH_K = 80
SPARSE_PREFETCH_K = 200
RETRIEVE_K = 150
RERANK_LIMIT = 15
TOP_K = 50
HTTP_TIMEOUT = 30.0
HTTP_MAX_RETRIES = 2
REQUIRED_ENV_VARS = ["EMBEDDINGS_DENSE_URL", "RERANKER_URL", "QDRANT_URL"]
logging.basicConfig(level=os.getenv("LOG_LEVEL", "INFO"))
logger = logging.getLogger("search-service")
def validate_required_env() -> None:
if bool(OPEN_API_LOGIN) != bool(OPEN_API_PASSWORD):
raise RuntimeError("OPEN_API_LOGIN and OPEN_API_PASSWORD must be set together")
if not API_KEY and not (OPEN_API_LOGIN and OPEN_API_PASSWORD):
raise RuntimeError("Either API_KEY or OPEN_API_LOGIN and OPEN_API_PASSWORD must be set")
missing = [name for name in REQUIRED_ENV_VARS if not os.getenv(name)]
if missing:
logger.error("Empty required env vars: %s", ", ".join(missing))
raise RuntimeError(f"Empty required env vars: {', '.join(missing)}")
def get_upstream_kwargs() -> dict[str, Any]:
headers = {"Content-Type": "application/json"}
kwargs: dict[str, Any] = {"headers": headers}
if OPEN_API_LOGIN and OPEN_API_PASSWORD:
kwargs["auth"] = (OPEN_API_LOGIN, OPEN_API_PASSWORD)
return kwargs
if API_KEY:
headers["Authorization"] = f"Bearer {API_KEY}"
return kwargs

View file

@ -1,4 +1,3 @@
import asyncio
import logging
import os
from contextlib import asynccontextmanager
@ -15,8 +14,9 @@ from qdrant_client import AsyncQdrantClient, models
EMBEDDINGS_DENSE_MODEL = "Qwen/Qwen3-Embedding-0.6B"
# Ваш сервис должен считывать эти переменные из окружения (env), так как проверяющая система управляет ими
HOST = os.getenv("HOST", "0.0.0.0")
PORT = int(os.getenv("PORT", "8000"))
PORT = int(os.getenv("PORT", "8003"))
API_KEY = os.getenv("API_KEY")
EMBEDDINGS_DENSE_URL = os.getenv("EMBEDDINGS_DENSE_URL")
@ -73,6 +73,7 @@ def get_upstream_request_kwargs() -> dict[str, Any]:
return kwargs
# Модель данных, которую мы предоставляем и рассчитываем получать от вас
class DateRange(BaseModel):
from_: str = Field(alias="from")
to: str
@ -125,9 +126,13 @@ class SparseVector(BaseModel):
values: list[float] = Field(default_factory=list)
class SparseEmbeddingResponse(BaseModel):
vectors: list[SparseVector]
# Метадата чанков в Qdrant'e, по которой вы можете фильтровать
class ChunkMetadata(BaseModel):
chat_name: str
chat_type: str
chat_type: str # channel, group, private, thread
chat_id: str
chat_sn: str
thread_sn: str | None = None
@ -162,14 +167,21 @@ async def lifespan(app: FastAPI):
app = FastAPI(title="Search Service", version="0.1.0", lifespan=lifespan)
DENSE_PREFETCH_K = 120
SPARSE_PREFETCH_K = 200
RETRIEVE_K = 150
RERANK_LIMIT = 35
KEYWORD_BOOST_EXTRA = 10
# Внутри шаблона dense и rerank берутся из внешних HTTP endpoint'ов,
# которые предоставляет проверяющая система.
# Текущий код ниже — минимальный пример search pipeline.
DENSE_PREFETCH_K = 30
SPARSE_PREFETCH_K = 40
RETRIEVE_K = 80
RERANK_LIMIT = 20
FINAL_TOP_K = 50
MAX_DENSE_QUERIES = 4
MAX_SPARSE_QUERIES = 3
async def embed_dense(client: httpx.AsyncClient, text: str) -> list[float]:
# Dense endpoint ожидает OpenAI-compatible body с input как списком строк.
response = await client.post(
EMBEDDINGS_DENSE_URL,
**get_upstream_request_kwargs(),
@ -179,31 +191,19 @@ async def embed_dense(client: httpx.AsyncClient, text: str) -> list[float]:
},
)
response.raise_for_status()
payload = DenseEmbeddingResponse.model_validate(response.json())
if not payload.data:
raise ValueError("Dense embedding response is empty")
return payload.data[0].embedding
async def embed_dense_batch(client: httpx.AsyncClient, texts: list[str]) -> list[list[float]]:
response = await client.post(
EMBEDDINGS_DENSE_URL,
**get_upstream_request_kwargs(),
json={
"model": os.getenv("EMBEDDINGS_DENSE_MODEL", EMBEDDINGS_DENSE_MODEL),
"input": texts,
},
)
response.raise_for_status()
payload = DenseEmbeddingResponse.model_validate(response.json())
payload.data.sort(key=lambda x: x.index)
return [item.embedding for item in payload.data]
def embed_sparse_sync(text: str) -> SparseVector:
async def embed_sparse(text: str) -> SparseVector:
vectors = list(get_sparse_model().embed([text]))
if not vectors:
raise ValueError("Sparse embedding response is empty")
item = vectors[0]
return SparseVector(
indices=[int(index) for index in item.indices.tolist()],
@ -211,93 +211,115 @@ def embed_sparse_sync(text: str) -> SparseVector:
)
def build_dense_query(question: Question) -> str:
q = question.search_text.strip() if question.search_text else question.text.strip()
return q
def unique_non_empty(values: list[str | None]) -> list[str]:
result: list[str] = []
seen: set[str] = set()
for value in values:
text = (value or "").strip()
if text and text not in seen:
seen.add(text)
result.append(text)
return result
def build_sparse_query(question: Question) -> str:
base = question.search_text.strip() if question.search_text else question.text.strip()
parts = [base]
if question.keywords:
parts.extend(question.keywords)
return " ".join(parts)
def question_entity_terms(question: Question) -> list[str]:
entities = question.entities
if entities is None:
return []
values: list[str | None] = []
values.extend(entities.people or [])
values.extend(entities.emails or [])
values.extend(entities.documents or [])
values.extend(entities.names or [])
values.extend(entities.links or [])
return unique_non_empty(values)
def _build_keyword_set(question: Question) -> list[str]:
tokens: list[str] = []
if question.keywords:
tokens.extend(kw.lower() for kw in question.keywords if kw)
if question.entities:
for field in (
question.entities.people,
question.entities.emails,
question.entities.documents,
question.entities.names,
question.entities.links,
):
tokens.extend(e.lower() for e in (field or []) if e)
return tokens
def build_query_texts(question: Question) -> tuple[str, list[str], list[str]]:
primary_query = (question.search_text or question.text).strip()
dense_queries = unique_non_empty(
[
primary_query,
*(question.variants or []),
*(question.hyde or []),
]
)[:MAX_DENSE_QUERIES]
keyword_query = " ".join(question.keywords or []).strip()
entity_query = " ".join(question_entity_terms(question)).strip()
sparse_queries = unique_non_empty(
[
keyword_query,
entity_query,
primary_query,
*(question.variants or []),
]
)[:MAX_SPARSE_QUERIES]
return primary_query, dense_queries, sparse_queries
def prefilter_for_rerank(
points: list[Any],
question: Question,
) -> tuple[list[Any], list[Any]]:
"""Select candidates for reranking: top by RRF + keyword-boosted stragglers."""
if not points:
return [], []
def build_search_filter(question: Question) -> models.Filter | None:
must_conditions: list[Any] = []
head = points[:RERANK_LIMIT]
tail = points[RERANK_LIMIT:]
if question.date_range and hasattr(models, "DatetimeRange"):
must_conditions.append(
models.FieldCondition(
key="metadata.start",
range=models.DatetimeRange(
gte=question.date_range.from_,
lte=question.date_range.to,
),
)
)
keywords = _build_keyword_set(question)
if not keywords or not tail:
return head, tail
extra: list[Any] = []
remaining_tail: list[Any] = []
for p in tail:
if len(extra) >= KEYWORD_BOOST_EXTRA:
remaining_tail.append(p)
continue
content = ((p.payload or {}).get("page_content") or "").lower()
if any(kw in content for kw in keywords):
extra.append(p)
else:
remaining_tail.append(p)
return head + extra, remaining_tail
return models.Filter(must=must_conditions) if must_conditions else None
async def qdrant_search(
client: AsyncQdrantClient,
dense_vectors: list[list[float]],
sparse_vector: SparseVector,
) -> list[Any] | None:
prefetch_list = []
for dv in dense_vectors:
prefetch_list.append(
sparse_vectors: list[SparseVector],
question_data: Question,
) -> Any | None:
search_filter = build_search_filter(question_data)
prefetch: list[models.Prefetch] = []
for dense_vector in dense_vectors:
prefetch.append(
models.Prefetch(
query=dv,
query=dense_vector,
using=QDRANT_DENSE_VECTOR_NAME,
limit=DENSE_PREFETCH_K,
filter=search_filter,
)
)
prefetch_list.append(
models.Prefetch(
query=models.SparseVector(
indices=sparse_vector.indices,
values=sparse_vector.values,
),
using=QDRANT_SPARSE_VECTOR_NAME,
limit=SPARSE_PREFETCH_K,
for sparse_vector in sparse_vectors:
if not sparse_vector.indices:
continue
prefetch.append(
models.Prefetch(
query=models.SparseVector(
indices=sparse_vector.indices,
values=sparse_vector.values,
),
using=QDRANT_SPARSE_VECTOR_NAME,
limit=SPARSE_PREFETCH_K,
filter=search_filter,
)
)
)
if not prefetch:
return None
response = await client.query_points(
collection_name=QDRANT_COLLECTION_NAME,
prefetch=prefetch_list,
prefetch=prefetch,
query=models.FusionQuery(fusion=models.Fusion.RRF),
limit=RETRIEVE_K,
with_payload=True,
@ -309,10 +331,37 @@ async def qdrant_search(
return response.points
def deduplicate_points(points: list[Any]) -> list[Any]:
unique_points: list[Any] = []
seen_ids: set[str] = set()
for point in points:
point_id = getattr(point, "id", None)
if point_id is None:
unique_points.append(point)
continue
point_id = str(point_id)
if point_id in seen_ids:
continue
seen_ids.add(point_id)
unique_points.append(point)
return unique_points
def extract_point_score(point: Any) -> float:
score = getattr(point, "score", 0.0)
if score is None:
return 0.0
return float(score)
def extract_message_ids(point: Any) -> list[str]:
payload = point.payload or {}
metadata = payload.get("metadata") or {}
message_ids = metadata.get("message_ids") or []
return [str(message_id) for message_id in message_ids]
@ -324,63 +373,95 @@ async def get_rerank_scores(
if not targets:
return []
for attempt in range(5):
try:
response = await client.post(
RERANKER_URL,
**get_upstream_request_kwargs(),
json={
"model": RERANKER_MODEL,
"encoding_format": "float",
"text_1": label,
"text_2": targets,
},
)
if response.status_code == 429:
wait = 2 ** attempt
logger.warning(f"Rerank 429, retry {attempt+1}/5 in {wait}s")
await asyncio.sleep(wait)
continue
response.raise_for_status()
payload = response.json()
data = payload.get("data") or []
return [float(sample["score"]) for sample in data]
except Exception as e:
logger.warning(f"Rerank error attempt {attempt+1}/5: {e}")
if attempt < 4:
await asyncio.sleep(2 ** attempt)
continue
logger.error("Rerank failed after 5 attempts, using fallback")
return []
# Rerank endpoint возвращает score для пары query -> candidate text.
response = await client.post(
RERANKER_URL,
**get_upstream_request_kwargs(),
json={
"model": RERANKER_MODEL,
"encoding_format": "float",
"text_1": label,
"text_2": targets,
},
)
response.raise_for_status()
logger.error("Rerank 429 after 5 retries, using fallback")
return []
payload = response.json()
data = payload.get("data") or []
return [float(sample["score"]) for sample in data]
async def rerank_points(
client: httpx.AsyncClient,
query: str,
points: list[Any],
) -> list[Any]:
if not points:
return []
targets = [point.payload.get("page_content") for point in points]
scores = await get_rerank_scores(client, query, targets)
) -> list[tuple[Any, float]]:
rerank_candidates = points[:RERANK_LIMIT]
tail_candidates = points[RERANK_LIMIT:]
rerank_targets = [
str((point.payload or {}).get("page_content") or "")
for point in rerank_candidates
]
if not scores or len(scores) != len(points):
logger.warning("Reranker unavailable or score mismatch, returning RRF order")
return points
try:
scores = await get_rerank_scores(client, query, rerank_targets)
except Exception:
logger.exception("Rerank failed, returning retrieval order")
return [(point, extract_point_score(point)) for point in points]
return [
point
for _, point in sorted(
zip(scores, points),
if len(scores) != len(rerank_candidates):
logger.warning(
"Rerank returned %d scores for %d candidates",
len(scores),
len(rerank_candidates),
)
return [(point, extract_point_score(point)) for point in points]
reranked_candidates = [
(point, float(score))
for score, point in sorted(
zip(scores, rerank_candidates),
key=lambda item: item[0],
reverse=True,
)
]
tail_with_scores = [
(point, extract_point_score(point))
for point in tail_candidates
]
return reranked_candidates + tail_with_scores
def aggregate_message_scores(
scored_points: list[tuple[Any, float]],
) -> tuple[dict[str, float], dict[str, int]]:
aggregated_scores: dict[str, float] = {}
first_seen_rank: dict[str, int] = {}
for rank, (point, point_score) in enumerate(scored_points):
point_message_ids = set(extract_message_ids(point))
for message_id in point_message_ids:
aggregated_scores[message_id] = aggregated_scores.get(message_id, 0.0) + point_score
first_seen_rank.setdefault(message_id, rank)
return aggregated_scores, first_seen_rank
def select_top_message_ids(
aggregated_scores: dict[str, float],
first_seen_rank: dict[str, int],
limit: int,
) -> list[str]:
sorted_items = sorted(
aggregated_scores.items(),
key=lambda item: (-item[1], first_seen_rank[item[0]]),
)
return [message_id for message_id, _ in sorted_items[:limit]]
# Ваш сервис должен имплементировать оба этих метода
@app.get("/health")
async def health() -> dict[str, str]:
return {"status": "ok"}
@ -388,65 +469,27 @@ async def health() -> dict[str, str]:
@app.post("/search", response_model=SearchAPIResponse)
async def search(payload: SearchAPIRequest) -> SearchAPIResponse:
question = payload.question
query = question.text.strip()
query, dense_queries, sparse_queries = build_query_texts(payload.question)
if not query:
raise HTTPException(status_code=400, detail="question.text is required")
raise HTTPException(status_code=400, detail="question.search_text or question.text is required")
client: httpx.AsyncClient = app.state.http
qdrant: AsyncQdrantClient = app.state.qdrant
dense_query = build_dense_query(question)
sparse_query = build_sparse_query(question)
dense_vectors = [await embed_dense(client, item) for item in dense_queries]
sparse_vectors = [await embed_sparse(item) for item in sparse_queries]
best_points = await qdrant_search(qdrant, dense_vectors, sparse_vectors, payload.question)
dense_task = embed_dense(client, dense_query)
sparse_task = asyncio.to_thread(lambda: embed_sparse_sync(sparse_query))
dense_vector, sparse_vector = await asyncio.gather(dense_task, sparse_task)
dense_vectors = [dense_vector]
extra_texts: list[str] = []
raw_text = question.text.strip()
if raw_text and raw_text != dense_query:
extra_texts.append(raw_text)
for v in (question.variants or []):
q_v = v.strip()
if q_v and q_v != dense_query and q_v not in extra_texts:
extra_texts.append(q_v)
for h in (question.hyde or []):
q_h = h.strip()
if q_h and q_h != dense_query and q_h not in extra_texts:
extra_texts.append(q_h)
extra_texts = extra_texts[:3]
if extra_texts:
try:
extra_vecs = await embed_dense_batch(client, extra_texts)
dense_vectors.extend(extra_vecs)
except Exception as e:
logger.warning(f"Extra dense embedding failed: {e}")
all_points = await qdrant_search(qdrant, dense_vectors, sparse_vector)
if all_points is None:
if not best_points:
return SearchAPIResponse(results=[])
all_points = list(all_points)
scored_points = await rerank_points(client, query, deduplicate_points(list(best_points)))
aggregated_scores, first_seen_rank = aggregate_message_scores(scored_points)
message_ids = select_top_message_ids(aggregated_scores, first_seen_rank, FINAL_TOP_K)
rerank_pool, rerank_tail = prefilter_for_rerank(all_points, question)
reranked = await rerank_points(client, query, rerank_pool)
final_points = reranked + rerank_tail
msg_score: dict[str, float] = {}
for rank, point in enumerate(reranked):
score = 1.0 / (rank + 1)
for mid in extract_message_ids(point):
msg_score[mid] = msg_score.get(mid, 0.0) + score
for rank, point in enumerate(rerank_tail):
score = 1.0 / (60 + rank + 1)
for mid in extract_message_ids(point):
msg_score[mid] = msg_score.get(mid, 0.0) + score
message_ids = sorted(msg_score, key=lambda m: msg_score[m], reverse=True)[:50]
return SearchAPIResponse(results=[SearchAPIItem(message_ids=message_ids)])
return SearchAPIResponse(
results=[SearchAPIItem(message_ids=message_ids)]
)
@app.exception_handler(Exception)
@ -466,7 +509,12 @@ async def exception_handler(request: Request, exc: Exception) -> JSONResponse:
def main() -> None:
import uvicorn
uvicorn.run("main:app", host=HOST, port=PORT, reload=False)
uvicorn.run(
"main:app",
host=HOST,
port=PORT,
reload=False,
)
if __name__ == "__main__":

View file

@ -1,109 +0,0 @@
import os
import re
from functools import lru_cache
import httpx
from fastembed import SparseTextEmbedding
from config import (
EMBEDDINGS_DENSE_MODEL,
EMBEDDINGS_DENSE_URL,
SPARSE_MODEL_NAME,
get_upstream_kwargs,
logger,
)
from schemas import DenseEmbeddingResponse, Question, SparseVector
@lru_cache(maxsize=1)
def get_sparse_model() -> SparseTextEmbedding:
logger.info("Loading local sparse model %s", SPARSE_MODEL_NAME)
return SparseTextEmbedding(model_name=SPARSE_MODEL_NAME)
async def embed_dense(client: httpx.AsyncClient, text: str) -> list[float]:
response = await client.post(
str(EMBEDDINGS_DENSE_URL),
**get_upstream_kwargs(),
json={
"model": os.getenv("EMBEDDINGS_DENSE_MODEL", EMBEDDINGS_DENSE_MODEL),
"input": [text],
},
)
response.raise_for_status()
payload = DenseEmbeddingResponse.model_validate(response.json())
if not payload.data:
raise ValueError("Dense embedding response is empty")
return payload.data[0].embedding
async def embed_dense_batch(client: httpx.AsyncClient, texts: list[str]) -> list[list[float]]:
"""Single request for multiple texts — avoids N parallel calls and rate limiting."""
response = await client.post(
str(EMBEDDINGS_DENSE_URL),
**get_upstream_kwargs(),
json={
"model": os.getenv("EMBEDDINGS_DENSE_MODEL", EMBEDDINGS_DENSE_MODEL),
"input": texts,
},
)
response.raise_for_status()
payload = DenseEmbeddingResponse.model_validate(response.json())
payload.data.sort(key=lambda x: x.index)
return [item.embedding for item in payload.data]
def embed_sparse(text: str) -> SparseVector:
vectors = list(get_sparse_model().embed([text]))
if not vectors:
raise ValueError("Sparse embedding response is empty")
item = vectors[0]
return SparseVector(
indices=[int(i) for i in item.indices.tolist()],
values=[float(v) for v in item.values.tolist()],
)
def _normalize_query(text: str) -> str:
return re.sub(r"\s+", " ", text).strip()
def build_primary_query(question: Question) -> str:
q = question.search_text.strip() if question.search_text else ""
if not q:
q = question.text.strip()
return _normalize_query(q)
def build_extra_dense_queries(question: Question) -> list[str]:
extras: list[str] = []
for v in question.variants or []:
q = _normalize_query(v)
if q:
extras.append(q)
for h in question.hyde or []:
q = _normalize_query(h)
if q:
extras.append(q)
return extras
def build_sparse_query(question: Question) -> str:
kws = question.keywords or []
if kws:
return " ".join(kws)
return build_primary_query(question)
def build_entity_tokens(question: Question) -> list[str]:
tokens: list[str] = []
if question.entities:
for field in (
question.entities.people,
question.entities.emails,
question.entities.documents,
question.entities.names,
question.entities.links,
):
tokens.extend(field or [])
return [t.strip() for t in tokens if t.strip()]

View file

@ -1,73 +0,0 @@
import asyncio
from typing import Any
import httpx
from config import RERANK_LIMIT, RERANKER_MODEL, RERANKER_URL, get_upstream_kwargs, logger
from retrieval import extract_page_content
async def _get_rerank_scores(
client: httpx.AsyncClient,
query: str,
targets: list[str],
) -> list[float]:
if not targets:
return []
for attempt in range(5):
try:
response = await client.post(
str(RERANKER_URL),
**get_upstream_kwargs(),
json={
"model": RERANKER_MODEL,
"encoding_format": "float",
"text_1": query,
"text_2": targets,
},
)
except Exception as exc:
if attempt < 4:
await asyncio.sleep(2 ** attempt)
continue
raise exc
if response.status_code == 429:
wait = 2 ** attempt
logger.warning("Rerank 429, retry %d/5 in %ds", attempt + 1, wait)
await asyncio.sleep(wait)
continue
response.raise_for_status()
data = response.json().get("data") or []
return [float(sample["score"]) for sample in data]
logger.error("Rerank 429 after all retries, falling back")
return []
async def rerank_points(
client: httpx.AsyncClient,
query: str,
points: list[Any],
) -> tuple[list[Any], list[Any]]:
if not points:
return [], []
head = points[:RERANK_LIMIT]
tail = points[RERANK_LIMIT:]
targets = [extract_page_content(p) for p in head]
try:
scores = await _get_rerank_scores(client, query, targets)
except Exception as exc:
logger.warning("Rerank failed, using retrieval order: %s", exc)
return head, tail
if len(scores) != len(head):
logger.warning("Rerank score count mismatch, using retrieval order")
return head, tail
reranked = [p for _, p in sorted(zip(scores, head), key=lambda x: x[0], reverse=True)]
return reranked, tail

View file

@ -1,119 +0,0 @@
from typing import Any
from qdrant_client import AsyncQdrantClient, models
from config import (
DENSE_PREFETCH_K,
QDRANT_COLLECTION_NAME,
QDRANT_DENSE_VECTOR_NAME,
QDRANT_SPARSE_VECTOR_NAME,
RETRIEVE_K,
SPARSE_PREFETCH_K,
logger,
)
from schemas import Question, SparseVector
def _build_filter(question: Question) -> models.Filter | None:
must_conditions: list[models.Condition] = []
if question.date_range:
try:
must_conditions.append(
models.FieldCondition(
key="metadata.end",
range=models.Range(gte=question.date_range.from_),
)
)
must_conditions.append(
models.FieldCondition(
key="metadata.start",
range=models.Range(lte=question.date_range.to),
)
)
except Exception as e:
logger.warning("Date filter failed: %s", e)
if question.asker:
must_conditions.append(
models.FieldCondition(
key="metadata.participants",
match=models.MatchValue(value=question.asker),
)
)
return models.Filter(must=must_conditions) if must_conditions else None
async def qdrant_search(
client: AsyncQdrantClient,
primary_dense: list[float],
extra_dense: list[list[float]],
sparse_vector: SparseVector,
question: Question,
) -> list[Any]:
search_filter = _build_filter(question)
prefetch: list[models.Prefetch] = []
# Primary dense
prefetch.append(
models.Prefetch(
query=primary_dense,
using=QDRANT_DENSE_VECTOR_NAME,
limit=DENSE_PREFETCH_K,
filter=search_filter,
)
)
# Extra dense (variants / hyde) - smaller budget per query
extra_k = max(10, DENSE_PREFETCH_K // max(1, len(extra_dense)))
for vec in extra_dense:
prefetch.append(
models.Prefetch(
query=vec,
using=QDRANT_DENSE_VECTOR_NAME,
limit=extra_k,
filter=search_filter,
)
)
# Sparse
prefetch.append(
models.Prefetch(
query=models.SparseVector(
indices=sparse_vector.indices,
values=sparse_vector.values,
),
using=QDRANT_SPARSE_VECTOR_NAME,
limit=SPARSE_PREFETCH_K,
filter=search_filter,
)
)
response = await client.query_points(
collection_name=QDRANT_COLLECTION_NAME,
prefetch=prefetch,
query=models.FusionQuery(fusion=models.Fusion.RRF),
limit=RETRIEVE_K,
with_payload=True,
)
if not response.points:
logger.debug("Qdrant returned 0 points")
return []
logger.debug("Qdrant returned %d points", len(response.points))
return list(response.points)
def extract_message_ids(point: Any) -> list[str]:
payload = point.payload or {}
metadata = payload.get("metadata") or {}
message_ids = metadata.get("message_ids") or []
return [str(mid) for mid in message_ids]
def extract_page_content(point: Any) -> str:
payload = point.payload or {}
return payload.get("page_content") or ""

View file

@ -1,72 +0,0 @@
from pydantic import BaseModel, Field
class DateRange(BaseModel):
from_: str = Field(alias="from")
to: str
class Entities(BaseModel):
people: list[str] | None = None
emails: list[str] | None = None
documents: list[str] | None = None
names: list[str] | None = None
links: list[str] | None = None
class Question(BaseModel):
text: str
asker: str = ""
asked_on: str = ""
variants: list[str] | None = None
hyde: list[str] | None = None
keywords: list[str] | None = None
entities: Entities | None = None
date_mentions: list[str] | None = None
date_range: DateRange | None = None
search_text: str = ""
class SearchAPIRequest(BaseModel):
question: Question
class SearchAPIItem(BaseModel):
message_ids: list[str]
class SearchAPIResponse(BaseModel):
results: list[SearchAPIItem]
class DenseEmbeddingItem(BaseModel):
index: int
embedding: list[float]
class DenseEmbeddingResponse(BaseModel):
data: list[DenseEmbeddingItem]
class SparseVector(BaseModel):
indices: list[int] = Field(default_factory=list)
values: list[float] = Field(default_factory=list)
class SparseEmbeddingResponse(BaseModel):
vectors: list[SparseVector]
class ChunkMetadata(BaseModel):
chat_name: str
chat_type: str
chat_id: str
chat_sn: str
thread_sn: str | None = None
message_ids: list[str]
start: str
end: str
participants: list[str] = Field(default_factory=list)
mentions: list[str] = Field(default_factory=list)
contains_forward: bool = False
contains_quote: bool = False

View file

@ -1,59 +0,0 @@
# Skill: Automating Task Tracking and Git Push with Codex
## Overview
This skill involves using **OpenAI Codex** to automate the tracking of task changes, logging updates, and automatically pushing changes to a Git repository. It combines **Codex's ability to generate code** with the power of **Git automation** to keep track of development tasks, log updates, and push them to version control.
## Objectives
1. Track the status of development tasks (e.g., status, assignee, comments).
2. Log all changes to tasks in a local directory (`.ai_update/`).
3. Automatically commit and push changes to a Git repository (GitHub, GitLab, etc.).
4. Provide clear and structured task progress reports.
## Components
- **Codex System Prompt**: Codex will act as the task manager, processing task updates and tracking changes.
- **Log Files**: Changes will be saved as text files in `.ai_update/`.
- **Git Integration**: Each task update will be automatically committed and pushed to a Git repository.
### 1. Codex System Prompt
Codex needs a **system prompt** that instructs it to track tasks, store updates in `.ai_update/`, and commit them to Git. Here is the **system prompt** for Codex:
```python
"""
You are a highly capable task manager for a development team working on a software project. Your job is to:
1. Track the tasks and changes in the project, including assigning tasks to team members and tracking progress.
2. Log every change made to the tasks, including updates to the status, comments, and assigned team members.
3. Ensure that no task is missed and that progress is clearly reported.
4. Store all updates, status changes, and task comments in a folder named `.ai_update/` on the local disk.
5. Automatically commit and push changes to the Git repository (on GitHub, GitLab, or other Git platforms) whenever updates are made.
6. Provide clear and detailed reports on what has been done so far and what tasks remain.
The folder `.ai_update/` will contain a log file where every change to a task is recorded, along with:
- Task name
- Updated status (if changed)
- Assignee (if changed)
- Any new comments (with timestamps)
- A summary of task progress
For each change:
1. Save the task update to the `.ai_update/` folder as a new log file.
2. Commit the new log file with a descriptive commit message, such as "Updated task [task_name] status to [status]".
3. Push the changes to the remote Git repository, ensuring that the changes are tracked properly.
4. If no change occurs in a task, do not log or push.
Ensure that no steps are skipped in the task completion process, and when a task is marked as "completed", ensure that all relevant information has been logged and pushed.
Provide the following feedback format for every task:
1. Task name
2. Current status
3. Comments and updates
4. Task progress summary
The `.ai_update/` folder will serve as a log of your progress. You should always commit the updates to Git in a way that shows a clear history of the changes.
Example of the log entry in `.ai_update/` folder:
Task: [task_name]
- Status: [status]
- Assignee: [assignee_name]
- Comment: [comment]
- Timestamp: [timestamp]
"""

View file

View file

@ -1,58 +0,0 @@
"""Unit tests for search/aggregation.py"""
import sys
import os
_SEARCH_DIR = os.path.join(os.path.dirname(__file__), "..", "search")
sys.path.insert(0, _SEARCH_DIR)
os.environ.setdefault("EMBEDDINGS_DENSE_URL", "http://localhost/embed")
os.environ.setdefault("RERANKER_URL", "http://localhost/rerank")
os.environ.setdefault("QDRANT_URL", "http://localhost:6333")
os.environ.setdefault("API_KEY", "test-key")
from aggregation import aggregate_message_ids
from config import TOP_K
def _point(message_ids: list[str]):
"""Fake qdrant point with payload."""
class FakePoint:
payload = {"metadata": {"message_ids": message_ids}}
return FakePoint()
class TestAggregateMessageIds:
def test_empty_inputs(self):
result = aggregate_message_ids([], [])
assert result == []
def test_basic_dedup(self):
head = [_point(["m1", "m2"]), _point(["m2", "m3"])]
result = aggregate_message_ids(head, [])
assert result.count("m2") == 1
def test_head_before_tail(self):
head = [_point(["head_msg"])]
tail = [_point(["tail_msg"])]
result = aggregate_message_ids(head, tail)
assert result.index("head_msg") < result.index("tail_msg")
def test_top_k_limit(self):
# Create enough points to exceed TOP_K
points = [_point([f"m{i}"]) for i in range(TOP_K + 20)]
result = aggregate_message_ids(points, [])
assert len(result) <= TOP_K
def test_cross_point_dedup(self):
head = [_point(["shared"]), _point(["shared", "unique"])]
result = aggregate_message_ids(head, [])
assert result.count("shared") == 1
assert "unique" in result
def test_tail_fills_after_head(self):
head = [_point(["h1"])]
tail = [_point(["t1"]), _point(["t2"])]
result = aggregate_message_ids(head, tail)
assert "h1" in result
assert "t1" in result
assert "t2" in result

View file

@ -1,125 +0,0 @@
"""Unit tests for index/chunking.py"""
import sys
import os
_INDEX_DIR = os.path.join(os.path.dirname(__file__), "..", "index")
sys.path.insert(0, _INDEX_DIR)
from chunking import build_chunks, _split_windows, WINDOW_MAX_MESSAGES, TIME_GAP_SECONDS
from cleaning import CleanedMessage
from index_schemas import Message
def _make_message(id: str, time: int, text: str = "hello", **kwargs) -> Message:
defaults = dict(
thread_sn=None,
sender_id="user@example.com",
file_snippets="",
parts=None,
mentions=None,
member_event=None,
is_system=False,
is_hidden=False,
is_forward=False,
is_quote=False,
)
defaults.update(kwargs)
return Message(id=id, time=time, text=text, **defaults)
def _make_cleaned(id: str, time: int, text: str = "hello") -> CleanedMessage:
return CleanedMessage(
id=id,
sender_id="user@x.com",
time=time,
thread_sn=None,
text=text,
parts=[],
mentions=[],
member_event_text="",
file_info=[],
is_system=False,
is_forward=False,
is_quote=False,
)
class TestBuildChunks:
def test_empty_new_messages(self):
result = build_chunks([], [])
assert result == []
def test_single_message(self):
msgs = [_make_message("m1", 1000000, text="A simple message")]
result = build_chunks([], msgs)
assert len(result) == 1
assert "m1" in result[0].message_ids
def test_message_ids_preserved(self):
msgs = [
_make_message("m1", 1000000, text="First"),
_make_message("m2", 1000100, text="Second"),
]
result = build_chunks([], msgs)
all_ids = [mid for chunk in result for mid in chunk.message_ids]
assert "m1" in all_ids
assert "m2" in all_ids
def test_different_content_fields(self):
msgs = [_make_message("m1", 1000000, text="test")]
result = build_chunks([], msgs)
chunk = result[0]
assert chunk.page_content
assert chunk.dense_content
assert chunk.sparse_content
def test_overlap_appears_in_chunk(self):
overlap = [_make_message("o1", 999000, text="overlap message")]
new_msgs = [_make_message("m1", 1000000, text="new message")]
result = build_chunks(overlap, new_msgs)
assert len(result) >= 1
# overlap ids should NOT be in message_ids (they're context only)
assert "o1" not in result[0].message_ids
assert "m1" in result[0].message_ids
def test_time_gap_splits_window(self):
msgs = [
_make_message("m1", 1000000, text="morning message"),
_make_message("m2", 1000000 + TIME_GAP_SECONDS + 1, text="evening message"),
]
result = build_chunks([], msgs)
# large time gap should create 2 chunks
assert len(result) == 2
def test_empty_messages_skipped(self):
msgs = [
_make_message("m1", 1000000, text=""),
_make_message("m2", 1000100, text="real content"),
]
result = build_chunks([], msgs)
all_ids = [mid for chunk in result for mid in chunk.message_ids]
assert "m1" not in all_ids
assert "m2" in all_ids
class TestSplitWindows:
def test_empty(self):
assert _split_windows([]) == []
def test_single(self):
msgs = [_make_cleaned("m1", 1000000)]
windows = _split_windows(msgs)
assert len(windows) == 1
def test_time_gap_splits(self):
msgs = [
_make_cleaned("m1", 1000000),
_make_cleaned("m2", 1000000 + TIME_GAP_SECONDS + 1),
]
windows = _split_windows(msgs)
assert len(windows) == 2
def test_max_messages_splits(self):
msgs = [_make_cleaned(f"m{i}", 1000000 + i * 10) for i in range(WINDOW_MAX_MESSAGES + 2)]
windows = _split_windows(msgs)
assert len(windows) >= 2

View file

@ -1,170 +0,0 @@
"""Unit tests for index/cleaning.py"""
import sys
import os
_INDEX_DIR = os.path.join(os.path.dirname(__file__), "..", "index")
sys.path.insert(0, _INDEX_DIR)
import pytest
from cleaning import (
normalize_unicode,
parse_file_snippets,
normalize_member_event,
normalize_part,
clean_message,
)
from index_schemas import Message
def _make_message(**kwargs) -> Message:
defaults = dict(
id="msg1",
thread_sn=None,
time=1000000,
text="",
sender_id="user@example.com",
file_snippets="",
parts=None,
mentions=None,
member_event=None,
is_system=False,
is_hidden=False,
is_forward=False,
is_quote=False,
)
defaults.update(kwargs)
return Message(**defaults)
class TestNormalizeUnicode:
def test_removes_zero_width(self):
assert "\u200b" not in normalize_unicode("hello\u200bworld")
assert "\u200c" not in normalize_unicode("a\u200cb")
assert "\ufeff" not in normalize_unicode("\ufefftext")
def test_collapses_whitespace(self):
result = normalize_unicode(" too many spaces ")
assert " " not in result
def test_normalizes_newlines(self):
result = normalize_unicode("line1\r\nline2\rline3")
assert "\r" not in result
def test_collapses_multiple_newlines(self):
result = normalize_unicode("a\n\n\n\nb")
assert "\n\n\n" not in result
def test_preserves_url(self):
url = "https://example.com/path?q=1&page=2"
assert url in normalize_unicode(url)
def test_preserves_email(self):
email = "user@corp.example"
assert email in normalize_unicode(email)
class TestParseFileSnippets:
def test_empty_string(self):
assert parse_file_snippets("") == []
def test_valid_list(self):
raw = '[{"name": "doc.pdf", "mime": "application/pdf"}]'
result = parse_file_snippets(raw)
assert len(result) == 1
assert result[0]["name"] == "doc.pdf"
def test_valid_dict(self):
raw = '{"name": "file.txt"}'
result = parse_file_snippets(raw)
assert len(result) == 1
def test_invalid_json(self):
assert parse_file_snippets("{broken json}") == []
def test_whitespace_only(self):
assert parse_file_snippets(" ") == []
class TestNormalizeMemberEvent:
def test_none_event(self):
assert normalize_member_event(None) == ""
def test_add_members(self):
event = {"type": "addMembers", "members": ["alice@example.com", "bob@example.com"]}
result = normalize_member_event(event)
assert "alice@example.com" in result
assert "added to chat" in result
def test_unknown_event(self):
event = {"type": "banUser", "member": "x@example.com"}
result = normalize_member_event(event)
assert "banUser" in result
class TestNormalizePart:
def test_text_part(self):
part = {"mediaType": "text", "text": "hello"}
result = normalize_part(part)
assert result["type"] == "text"
assert result["text"] == "hello"
def test_quote_part(self):
part = {"mediaType": "quote", "sn": "alice@x.com", "text": "original text"}
result = normalize_part(part)
assert result["type"] == "quote"
assert "alice@x.com" in result["text"]
assert "original text" in result["text"]
def test_forward_part(self):
part = {"mediaType": "forward", "sn": "channel@x.com", "text": "forwarded"}
result = normalize_part(part)
assert result["type"] == "forward"
assert "forwarded" in result["text"]
def test_unknown_media_type(self):
part = {"mediaType": "sticker", "text": ""}
result = normalize_part(part)
assert "sticker" in result["type"]
class TestCleanMessage:
def test_empty_message_is_empty(self):
msg = _make_message()
cleaned = clean_message(msg)
assert cleaned.is_empty
def test_text_message(self):
msg = _make_message(text="Hello world")
cleaned = clean_message(msg)
assert not cleaned.is_empty
assert cleaned.text == "Hello world"
def test_parts_extracted(self):
msg = _make_message(
parts=[{"mediaType": "text", "text": "from parts"}]
)
cleaned = clean_message(msg)
assert not cleaned.is_empty
assert any("from parts" in p["text"] for p in cleaned.parts)
def test_member_event_extracted(self):
msg = _make_message(
is_system=True,
member_event={"type": "addMembers", "members": ["u@x.com"]},
)
cleaned = clean_message(msg)
assert not cleaned.is_empty
assert "u@x.com" in cleaned.member_event_text
def test_file_snippets_parsed(self):
msg = _make_message(
file_snippets='[{"name": "report.pdf", "mime": "application/pdf"}]'
)
cleaned = clean_message(msg)
assert not cleaned.is_empty
assert cleaned.file_info[0]["name"] == "report.pdf"
def test_zero_width_stripped_from_text(self):
msg = _make_message(text="hello\u200bworld")
cleaned = clean_message(msg)
assert "\u200b" not in cleaned.text

View file

@ -1,118 +0,0 @@
"""Unit tests for search/query_builder.py (pure logic only, no HTTP)"""
import sys
import os
_SEARCH_DIR = os.path.join(os.path.dirname(__file__), "..", "search")
sys.path.insert(0, _SEARCH_DIR)
# Stub env vars before importing search modules
os.environ.setdefault("EMBEDDINGS_DENSE_URL", "http://localhost/embed")
os.environ.setdefault("RERANKER_URL", "http://localhost/rerank")
os.environ.setdefault("QDRANT_URL", "http://localhost:6333")
os.environ.setdefault("API_KEY", "test-key")
from schemas import Entities, Question
from query_builder import (
build_primary_query,
build_extra_dense_queries,
build_sparse_query,
build_entity_tokens,
)
def _q(**kwargs) -> Question:
defaults = dict(text="default question")
defaults.update(kwargs)
return Question(**defaults)
class TestBuildPrimaryQuery:
def test_uses_search_text_over_text(self):
q = _q(text="original", search_text="refined query")
assert build_primary_query(q) == "refined query"
def test_fallback_to_text(self):
q = _q(text="fallback text", search_text="")
assert build_primary_query(q) == "fallback text"
def test_strips_whitespace(self):
q = _q(text=" trimmed ")
assert build_primary_query(q) == "trimmed"
def test_collapses_internal_spaces(self):
q = _q(text="too many spaces")
result = build_primary_query(q)
assert " " not in result
class TestBuildExtraDenseQueries:
def test_no_extras_when_none(self):
q = _q(text="q")
assert build_extra_dense_queries(q) == []
def test_includes_variants(self):
q = _q(text="q", variants=["var1", "var2"])
extras = build_extra_dense_queries(q)
assert "var1" in extras
assert "var2" in extras
def test_includes_hyde(self):
q = _q(text="q", hyde=["hypothetical answer"])
extras = build_extra_dense_queries(q)
assert "hypothetical answer" in extras
def test_skips_empty_strings(self):
q = _q(text="q", variants=["", " ", "valid"])
extras = build_extra_dense_queries(q)
assert "" not in extras
assert " " not in extras
assert "valid" in extras
class TestBuildSparseQuery:
def test_uses_keywords_when_present(self):
q = _q(text="question", keywords=["go", "golang", "performance"])
result = build_sparse_query(q)
assert "go" in result
assert "golang" in result
def test_fallback_to_primary_when_no_keywords(self):
q = _q(text="fallback question", search_text="refined")
result = build_sparse_query(q)
assert result == "refined"
def test_empty_keywords_fallback(self):
q = _q(text="my question", keywords=[])
result = build_sparse_query(q)
assert result == "my question"
class TestBuildEntityTokens:
def test_no_entities(self):
q = _q(text="q")
assert build_entity_tokens(q) == []
def test_people_extracted(self):
q = _q(text="q", entities=Entities(people=["Alice", "Bob"]))
tokens = build_entity_tokens(q)
assert "Alice" in tokens
assert "Bob" in tokens
def test_all_entity_fields(self):
q = _q(
text="q",
entities=Entities(
people=["Alice"],
emails=["alice@corp.com"],
documents=["report.pdf"],
names=["Project X"],
links=["https://example.com"],
),
)
tokens = build_entity_tokens(q)
assert len(tokens) == 5
def test_strips_whitespace(self):
q = _q(text="q", entities=Entities(people=[" Alice "]))
tokens = build_entity_tokens(q)
assert "Alice" in tokens

View file

@ -1,110 +0,0 @@
"""Unit tests for index/rendering.py"""
import sys
import os
_INDEX_DIR = os.path.join(os.path.dirname(__file__), "..", "index")
sys.path.insert(0, _INDEX_DIR)
from cleaning import CleanedMessage
from rendering import render_page_content, render_dense_content, render_sparse_content
def _make_cleaned(**kwargs) -> CleanedMessage:
defaults = dict(
id="msg1",
sender_id="alice@example.com",
time=1700000000,
thread_sn=None,
text="",
parts=[],
mentions=[],
member_event_text="",
file_info=[],
is_system=False,
is_forward=False,
is_quote=False,
)
defaults.update(kwargs)
return CleanedMessage(**defaults)
class TestRenderPageContent:
def test_text_message(self):
msg = _make_cleaned(text="Hello world")
result = render_page_content(msg)
assert "alice@example.com" in result
assert "Hello world" in result
def test_quote_part(self):
msg = _make_cleaned(
parts=[{"type": "quote", "text": "[quote from bob]: original"}]
)
result = render_page_content(msg)
assert "[quote from bob]" in result
def test_forward_part(self):
msg = _make_cleaned(
parts=[{"type": "forward", "text": "[forwarded from channel]: content"}],
is_forward=True,
)
result = render_page_content(msg)
assert "forwarded" in result
def test_member_event(self):
msg = _make_cleaned(
member_event_text="[system: alice@x.com added to chat]",
is_system=True,
)
result = render_page_content(msg)
assert "added to chat" in result
def test_file_attachment(self):
msg = _make_cleaned(
file_info=[{"name": "report.pdf", "mime": "application/pdf", "url": "", "date": ""}]
)
result = render_page_content(msg)
assert "report.pdf" in result
class TestRenderDenseContent:
def test_includes_timestamp(self):
msg = _make_cleaned(text="test", time=1700000000)
result = render_dense_content(msg)
assert "2023-" in result # UTC date
def test_includes_sender(self):
msg = _make_cleaned(text="hi", sender_id="bob@corp.com")
result = render_dense_content(msg)
assert "sender:bob@corp.com" in result
def test_forward_marker(self):
msg = _make_cleaned(is_forward=True, parts=[{"type": "forward", "text": "[forwarded]: x"}])
result = render_dense_content(msg)
assert "type:forward" in result
def test_mentions_included(self):
msg = _make_cleaned(
text="hey",
mentions=["charlie@corp.com"],
)
result = render_dense_content(msg)
assert "charlie@corp.com" in result
class TestRenderSparseContent:
def test_includes_sender(self):
msg = _make_cleaned(text="hello")
result = render_sparse_content(msg)
assert "alice@example.com" in result
def test_includes_mentions(self):
msg = _make_cleaned(text="ping", mentions=["dave@corp.com"])
result = render_sparse_content(msg)
assert "dave@corp.com" in result
def test_includes_filename(self):
msg = _make_cleaned(
file_info=[{"name": "budget.xlsx", "mime": "", "url": "", "date": ""}]
)
result = render_sparse_content(msg)
assert "budget.xlsx" in result