Compare commits
| Author | SHA256 | Date | |
|---|---|---|---|
|
|
5df69d73c7 | ||
|
|
4704f8485b | ||
|
|
9a7019a2ce | ||
|
|
07b1623c5e | ||
|
|
681437acb1 | ||
|
|
a183a586cb | ||
|
|
af62b51ef5 | ||
|
|
f640451257 | ||
|
|
a56252c361 | ||
|
|
38141ea5c7 | ||
|
|
d24d1c178b | ||
|
|
93673e5786 | ||
|
|
2f3b1571e5 | ||
|
|
65703f0fc4 | ||
|
|
4e674af4ed | ||
|
|
f9c2d7c63a | ||
|
|
968d85d738 | ||
|
|
d230b61c0b | ||
|
|
59047a8d0e | ||
|
|
81493ac4be | ||
|
|
cf33f9a3bd | ||
|
|
d2f65b169c | ||
|
|
fa1f7358b7 | ||
|
|
8208bb0b9d | ||
|
|
daa516babe | ||
|
|
6c2a9836eb | ||
|
|
bd3e6a56ca | ||
|
|
ea2ee1711c | ||
|
|
2890cdd457 | ||
|
|
e4dfb43b5e | ||
|
|
febdbc059e | ||
|
|
3926d561c0 | ||
|
|
d6bc883520 | ||
|
|
b2d20671cd | ||
|
|
a72e2e7014 | ||
|
|
cf6ca1d96b | ||
|
|
9d1e49949a | ||
|
|
ea3b34e11e | ||
|
|
09dc55db0e | ||
|
|
85133db483 | ||
|
|
6fc7d75e17 | ||
|
|
be321813fa | ||
|
|
d2c0ecdf62 | ||
|
|
266a74ef79 | ||
|
|
18d27c0146 | ||
|
|
8c4997ee2e | ||
|
|
127c561a41 |
@@ -13,6 +13,7 @@ build/
|
||||
.kotlin
|
||||
|
||||
### IntelliJ IDEA ###
|
||||
.idea/
|
||||
.idea/modules.xml
|
||||
.idea/jarRepositories.xml
|
||||
.idea/compiler.xml
|
||||
@@ -102,10 +103,12 @@ ESP32/**/*.d
|
||||
ESP32/**/*.a
|
||||
|
||||
# Полные серверные бэкапы (тяжёлые архивы, не коммитим)
|
||||
server-backup/archive/**
|
||||
!server-backup/archive/.gitkeep
|
||||
deploy/backup/archive/**
|
||||
!deploy/backup/archive/.gitkeep
|
||||
|
||||
# Локальная дев-обвязка Claude (дев-сервер shine-UI, сессии, планы) — не коммитим
|
||||
# Локальная дев-обвязка AI-агентов (сессии, планы, настройки) — не коммитим
|
||||
.agents/
|
||||
.codex/
|
||||
.claude/
|
||||
# Рабочие бэкапы/превью-ассеты UI — не для репозитория
|
||||
*.bak.png
|
||||
|
||||
@@ -1,8 +0,0 @@
|
||||
# Default ignored files
|
||||
/shelf/
|
||||
/workspace.xml
|
||||
# Editor-based HTTP Client requests
|
||||
/httpRequests/
|
||||
# Datasource local storage ignored files
|
||||
/dataSources/
|
||||
/dataSources.local.xml
|
||||
@@ -1 +0,0 @@
|
||||
shine-server-server
|
||||
@@ -1,10 +0,0 @@
|
||||
<component name="ArtifactManager">
|
||||
<artifact type="jar" build-on-make="true" name="server:jar">
|
||||
<output-path>$PROJECT_DIR$/out/artifacts/server_jar</output-path>
|
||||
<root id="archive" name="server.jar">
|
||||
<element id="directory" name="META-INF">
|
||||
<element id="file-copy" path="$PROJECT_DIR$/META-INF/MANIFEST.MF" />
|
||||
</element>
|
||||
</root>
|
||||
</artifact>
|
||||
</component>
|
||||
@@ -1,25 +0,0 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<project version="4">
|
||||
<component name="GradleMigrationSettings" migrationVersion="1" />
|
||||
<component name="GradleSettings">
|
||||
<option name="linkedExternalProjectsSettings">
|
||||
<GradleProjectSettings>
|
||||
<option name="externalProjectPath" value="$PROJECT_DIR$" />
|
||||
<option name="gradleHome" value="" />
|
||||
<option name="modules">
|
||||
<set>
|
||||
<option value="$PROJECT_DIR$" />
|
||||
<option value="$PROJECT_DIR$/shine-server-blockchain" />
|
||||
<option value="$PROJECT_DIR$/shine-server-config" />
|
||||
<option value="$PROJECT_DIR$/shine-server-crypto" />
|
||||
<option value="$PROJECT_DIR$/shine-server-db" />
|
||||
<option value="$PROJECT_DIR$/shine-server-geo" />
|
||||
<option value="$PROJECT_DIR$/shine-server-log" />
|
||||
<option value="$PROJECT_DIR$/shine-server-net-protocol" />
|
||||
<option value="$PROJECT_DIR$/shine-server-net-server" />
|
||||
</set>
|
||||
</option>
|
||||
</GradleProjectSettings>
|
||||
</option>
|
||||
</component>
|
||||
</project>
|
||||
@@ -1,10 +0,0 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<project version="4">
|
||||
<component name="ExternalStorageConfigurationManager" enabled="true" />
|
||||
<component name="FrameworkDetectionExcludesConfiguration">
|
||||
<file type="web" url="file://$PROJECT_DIR$" />
|
||||
</component>
|
||||
<component name="ProjectRootManager" version="2" languageLevel="JDK_17" default="true" project-jdk-name="17 (2)" project-jdk-type="JavaSDK">
|
||||
<output url="file://$PROJECT_DIR$/out" />
|
||||
</component>
|
||||
</project>
|
||||
@@ -1,6 +0,0 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<project version="4">
|
||||
<component name="VcsDirectoryMappings">
|
||||
<mapping directory="" vcs="Git" />
|
||||
</component>
|
||||
</project>
|
||||
@@ -14,14 +14,15 @@
|
||||
- Веб-панель администратора сервера (управление Solana PDA сервера) находится в `shine-UI/`:
|
||||
- точка входа `shine-UI/server-ui.html`;
|
||||
- остальные файлы серверного UI — в `shine-UI/server-ui/`.
|
||||
- Локальный Telegram-бот агента-кодера находится в папке `SHiNE-agent-bot-coder/` и не является кодом основного серверного приложения.
|
||||
- Локальный Telegram-бот агента-кодера живёт рядом с репозиторием продукта, обычно в `../SHiNE-agent-bot-coder/`, и не входит в публичный код основного приложения.
|
||||
- Solana/Anchor-модуль находится в папке `shine-solana/shine/` и ведётся отдельно от основного server/UI деплоя.
|
||||
|
||||
## Сервис агента-кодера
|
||||
- В проекте есть локальный Telegram-бот-сервис агента-кодера в папке `SHiNE-agent-bot-coder/`.
|
||||
- Локальный Telegram-бот-сервис агента-кодера находится вне этого git-репозитория, обычно в `../SHiNE-agent-bot-coder/`.
|
||||
- Сервис принимает сообщения из Telegram, ведёт историю диалога, ставит задачи в очередь и вызывает Codex CLI для обработки запросов по проекту.
|
||||
- Автоматически читаемые инструкции для Codex внутри сервиса держать в `SHiNE-agent-bot-coder/AGENTS.md`.
|
||||
- Подробные служебные правила Telegram-обработчика, его очередь, история, systemd-запуск и особенности ответов описывать в `SHiNE-agent-bot-coder/AGENT.md`.
|
||||
- Рабочая папка Codex для сервиса должна указывать на этот продуктовый репозиторий: `CODEX_WORKDIR=/home/ai/work/SHiNE/SHiNE-server-sha256/SHiNE-product`.
|
||||
- Автоматически читаемые инструкции для Codex внутри сервиса держать в `../SHiNE-agent-bot-coder/AGENTS.md`.
|
||||
- Подробные служебные правила Telegram-обработчика, его очередь, история, systemd-запуск и особенности ответов описывать в `../SHiNE-agent-bot-coder/AGENT.md`.
|
||||
- Если в сообщениях пользователя встречается «агент MD» или похожая формулировка про файл инструкций Codex, считать, что имеется в виду автоматически читаемый `AGENTS.md`.
|
||||
|
||||
## ESP32 UI homeserver
|
||||
@@ -36,70 +37,62 @@
|
||||
- Модуль логически связан с SHiNE, но не должен автоматически подключаться к сборке или деплою основного сервера без отдельного решения.
|
||||
- В Solana-модуле действуют локальные инструкции `shine-solana/shine/AGENTS.md`; при изменениях внутри модуля сначала читать их.
|
||||
- В git добавлять исходники, lock-файлы, настройки проекта и документацию Solana-модуля, но не добавлять локальные ключи, `.git`, `.idea`, `.gradle`, `target`, `node_modules`, `test-ledger`, логи, временные run-отчёты и `.env`-конфиги.
|
||||
- Для Solana deploy/push использовать правила из локального `shine-solana/shine/AGENTS.md`; не смешивать deploy Solana-модуля с `deployServer`/`deployUI` основного проекта.
|
||||
- Для Solana deploy/push использовать правила из локального `shine-solana/shine/AGENTS.md`; не смешивать deploy Solana-модуля со скриптами основного server/UI deploy из `deploy/scripts/`.
|
||||
- Для регистрации пользователей в Solana (программа `shine_users`) единая актуальная инструкция по деплою/инициализации, адресам программ, и куда их прописывать в UI/сервере находится в:
|
||||
- `Dev_Docs/Инициализация_Solana_регистрации/README.md`
|
||||
- `docs/Инициализация_Solana_регистрации/README.md`
|
||||
- Этот файл считать основной справкой (single source of truth) по деплою и первичной инициализации Solana-регистрации в текущем проекте.
|
||||
- Актуальная архитектурная справка по устройству Solana-программ, PDA-счетам, ролям DAO и движению средств находится в:
|
||||
- `Dev_Docs/Solana_Architecture/README.md`
|
||||
- `docs/Solana_Architecture/README.md`
|
||||
- Документ формата пользовательской PDA-записи `shine_users` находится в:
|
||||
- `shine-solana/shine/doc/formats/shine-user-pda-format-v.1.0.md`
|
||||
|
||||
## Документация блокчейна
|
||||
- Актуальная документация по форматам блокчейна находится в `Dev_Docs/Blockchain/README.md`.
|
||||
- Актуальная документация по форматам блокчейна находится в `docs/Blockchain/README.md`.
|
||||
- Это точка входа (оглавление), рядом расположены детальные файлы по форматам, типам каналов и командным сообщениям.
|
||||
- При любом изменении кода, связанного с блокчейном (формат блока, типы каналов, правила чтения/записи, команды), обязательно обновлять соответствующие документы в `Dev_Docs/Blockchain/`.
|
||||
- Дополнительно обязательно вести `Dev_Docs/Blockchain/CHANGELOG.md`: дописывать изменения построчно с указанием даты/времени и хэша коммита, после которого внесено изменение.
|
||||
- При любом изменении кода, связанного с блокчейном (формат блока, типы каналов, правила чтения/записи, команды), обязательно обновлять соответствующие документы в `docs/Blockchain/`.
|
||||
- Дополнительно обязательно вести `docs/Blockchain/CHANGELOG.md`: дописывать изменения построчно с указанием даты/времени и хэша коммита, после которого внесено изменение.
|
||||
- Перед любым изменением формата блокчейна обязательно заранее предупреждать пользователя, что формат будет изменён.
|
||||
- Изменять формат блокчейна можно только после явного подтверждения пользователя (без подтверждения формат не менять).
|
||||
- Добавление любых данных в блокчейн выполнять только через операцию `AddBlock`.
|
||||
- Перед каждым `AddBlock` обязательно проверять/актуализировать текущее состояние вершины блокчейна (`last global number/hash`) и использовать его при формировании блока.
|
||||
|
||||
## Документация личных сообщений (DM)
|
||||
- Актуальная документация по логике личных сообщений находится в `Dev_Docs/Personal_Messages/Протокол_DM_v1.md`.
|
||||
- Точный байтовый формат DM находится в `Dev_Docs/Personal_Messages/Формат_DM_v1.md`.
|
||||
- Актуальная документация по логике личных сообщений находится в `docs/Personal_Messages/Протокол_DM_v1.md`.
|
||||
- Точный байтовый формат DM находится в `docs/Personal_Messages/Формат_DM_v1.md`.
|
||||
- При любом изменении кода, связанного с личными сообщениями (формат подписанного DM-блока, типы DM-сообщений, правила доставки/ACK/read-receipt, роутинг по сессиям, UI-логика чатов), обязательно обновлять оба документа:
|
||||
- `Dev_Docs/Personal_Messages/Протокол_DM_v1.md`
|
||||
- `Dev_Docs/Personal_Messages/Формат_DM_v1.md`
|
||||
- `docs/Personal_Messages/Протокол_DM_v1.md`
|
||||
- `docs/Personal_Messages/Формат_DM_v1.md`
|
||||
- Логика личных сообщений в коде должна всегда соответствовать этим документам.
|
||||
- Документ по личным сообщениям обязан поддерживаться в актуальном состоянии.
|
||||
|
||||
## Документация API сервера
|
||||
- Актуальная документация по публичному JSON/WebSocket API сервера находится в `Dev_Docs/API/`.
|
||||
- При любом изменении серверного API/эндпоинтов/операций `op` обязательно обновлять соответствующие документы в `Dev_Docs/API/`.
|
||||
- Актуальная документация по публичному JSON/WebSocket API сервера находится в `docs/API/`.
|
||||
- При любом изменении серверного API/эндпоинтов/операций `op` обязательно обновлять соответствующие документы в `docs/API/`.
|
||||
- Перед изменением самого серверного API обязательно явно предупредить пользователя, какие операции, поля запросов/ответов или коды ошибок будут изменены, и запросить отдельное подтверждение.
|
||||
- Без явного подтверждения пользователя формат серверного API не менять; допускается только приведение документации в соответствие уже существующему коду.
|
||||
- Если добавляется новая операция `op`, нужно обновить общий список операций в `Dev_Docs/API/09_Operations_Index.md` или создать его, если файла ещё нет.
|
||||
- Если добавляется новая операция `op`, нужно обновить общий список операций в `docs/API/09_Operations_Index.md` или создать его, если файла ещё нет.
|
||||
|
||||
## Документация Figma
|
||||
- Актуальная документация по переносу экранов SHiNE в Figma и обратному переносу из Figma в код находится в `Dev_Docs/Figma/`.
|
||||
- Точка входа: `Dev_Docs/Figma/README.md`.
|
||||
- Подробный рабочий регламент: `Dev_Docs/Figma/TRANSFER_UI_SCREENS.md`.
|
||||
- Актуальная документация по переносу экранов SHiNE в Figma и обратному переносу из Figma в код находится в `docs/Figma/`.
|
||||
- Точка входа: `docs/Figma/README.md`.
|
||||
- Подробный рабочий регламент: `docs/Figma/TRANSFER_UI_SCREENS.md`.
|
||||
- Для экранов регистрации, входа и других чувствительных UI-flow по умолчанию переносить экраны в Figma по одному, а не пачкой, если пользователь отдельно не подтвердил иной способ.
|
||||
|
||||
## Версионирование
|
||||
- Единый файл версий проекта: `VERSION.properties` (в корне репозитория).
|
||||
- Перед каждым новым коммитом обязательно увеличивать версии в `VERSION.properties`:
|
||||
- `client.version` — версия клиентского UI.
|
||||
- `server.version` — версия серверной части.
|
||||
- Базовое правило инкремента: `+1` по последнему числовому сегменту (patch), если не оговорено иное.
|
||||
- Обычные коммиты делать стандартным `git commit`; переменная `$GITEA_TOKEN` для коммитов не нужна и не используется.
|
||||
- Все правила по коммитам, merge в `main`, `git push` и обновлению `VERSION.properties` находятся в `COMMIT_AND_VERSION_RULES.md`.
|
||||
- Этот файл считать единым источником истины по правилам версионирования и коммитов для данного репозитория.
|
||||
|
||||
## Deploy
|
||||
- Все документы и заметки по деплою хранить в папке `Dev_Docs/deploy/`.
|
||||
- Production-хост SHiNE: `player@shineup.me` (`178.208.64.62`).
|
||||
- Основной test-хост SHiNE: `player@193.8.215.70` (`t.shineup.me`).
|
||||
- Резервный test-хост SHiNE: `player@93.170.12.154` (`test.shineup.me`).
|
||||
- Базовый путь на сервере для SHiNE: `/home/player` (проекты SHiNE размещать в `/home/player/SHiNE/...`).
|
||||
- Все документы, инструкции, backup-правила и скрипты деплоя хранить в папке `deploy/`.
|
||||
- Подробные правила для агента по деплою находятся в `deploy/AGENTS.md`; перед любым deploy читать этот файл.
|
||||
- Production-серверы SHiNE: `shineup.me` и `server2.shineup.me`.
|
||||
- Тестовые/devnet серверы SHiNE: `t1.shineup.me`, `t2.shineup.me`, `t3.shineup.me`, `t4.shineup.me`.
|
||||
- В deploy-документах и скриптах использовать домены, а не IP.
|
||||
- По возможности все справки, комментарии и примечания в конфигах/документах писать на русском языке.
|
||||
- Для операций `git push` при необходимости использовать токен из переменной окружения `$GITEA_TOKEN`.
|
||||
- Любые изменения и любой деплой на production `shineup.me` выполнять только после отдельного явного подтверждения пользователя.
|
||||
- Если пользователь пишет просто `задеплой` без уточнения production/test, по умолчанию деплоить на `t.shineup.me`.
|
||||
- Default server deploy: `./gradlew deployServer` или `./gradlew deployServerTest2`.
|
||||
- Default UI deploy: `./gradlew deployUI` или `./gradlew deployUITest2`.
|
||||
- Production server deploy: `./gradlew deployServerProduction`.
|
||||
- Production UI deploy: `./gradlew deployUIProduction`.
|
||||
- Резервный test deploy на `test.shineup.me`: `./gradlew deployServerTest` и `./gradlew deployUITest`, но пока их не использовать без отдельной причины.
|
||||
- Любые изменения и любой деплой на production (`shineup.me` и `server2.shineup.me`) выполнять только после отдельного явного подтверждения пользователя.
|
||||
- Перед production deploy обязательно проверить/обновить бэкап в `deploy/backup/archive/`.
|
||||
- Если пользователь пишет просто `задеплой` без уточнения production/test, уточнить целевой контур; не выбирать production автоматически.
|
||||
- Deploy выполнять shell-скриптами из `deploy/scripts/`; Gradle deploy-задачи не использовать.
|
||||
- Для локального запуска использовать `./gradlew startLocal` (или `startLocalWithBuild`).
|
||||
- Сначала предлагать локальную проверку, а деплой на сервер выполнять по запросу пользователя.
|
||||
- Для временной бесплатной загрузки аватаров в Arweave секретный JWK нельзя хранить в git и нельзя прописывать в репозиторный `application.properties`.
|
||||
@@ -125,26 +118,13 @@
|
||||
- `unknown_error`
|
||||
- В этих записях искать поля `reason`, `failureStage`, `pcConnectionState`, `pcIceConnectionState`, `routeLabel`, `configuredTurnHosts*`, `reachableTurnHosts*`.
|
||||
|
||||
## Недопроверенные фичи (обязательно)
|
||||
- Папка для учёта недопроверенных фич: `Dev_Docs/Pending_Features/`.
|
||||
- По каждой новой доработке, которая требует ручной проверки, добавлять отдельный markdown-файл в `Dev_Docs/Pending_Features/`.
|
||||
- Рекомендуемый формат имени файла: `YYYY-MM-DD_HHMM_<short-feature-name>.md`.
|
||||
- Имена новых файлов и краткие описания фич по возможности писать на русском языке.
|
||||
- Внутри файла обязательно указывать:
|
||||
- краткое описание фичи;
|
||||
- что именно проверять;
|
||||
- ожидаемый результат;
|
||||
- статус (например: `pending`, `in_progress`, `done`).
|
||||
- После подтверждения, что фича проверена и работает корректно, соответствующий файл удалять.
|
||||
- В `Dev_Docs/Pending_Features/README.md` вести краткий регламент и поддерживать актуальность.
|
||||
|
||||
## Будущие фичи / TODO
|
||||
- Папка для задач, сознательно отложенных на будущее: `TODO/`.
|
||||
- Точка входа по планам: `TODO/README.md`.
|
||||
- Внутри планы разделены по горизонтам: `near/`, `medium/`, `far/` и тематическим подпапкам.
|
||||
- Если пользователь спрашивает, какие есть планы или что можно продолжить, сначала читать `TODO/README.md`, затем при необходимости конкретные файлы из подпапок.
|
||||
- Файлы из этой папки не считать активными задачами и не начинать реализацию без явной просьбы пользователя.
|
||||
- Старую папку `Dev_Docs/Future_Features/` считать выведенной из использования и больше не использовать для новых записей.
|
||||
- Старую папку `docs/Future_Features/` считать выведенной из использования и больше не использовать для новых записей.
|
||||
- Если часть кода временно отключена или закомментирована, либо удалена как временная заглушка, в TODO-файле подробно описывать:
|
||||
- какие файлы и участки отключены;
|
||||
- что осталось в коде как заготовка;
|
||||
|
||||
@@ -5,7 +5,7 @@
|
||||
@shine-UI/AGENTS.md
|
||||
|
||||
## Справка по подпроектам
|
||||
- При работе внутри `SHiNE-agent-bot-coder/` — читать `SHiNE-agent-bot-coder/AGENTS.md` и `SHiNE-agent-bot-coder/AGENT.md`.
|
||||
- При работе с локальным агентом-кодером — читать внешние файлы `../SHiNE-agent-bot-coder/AGENTS.md` и `../SHiNE-agent-bot-coder/AGENT.md`.
|
||||
- При работе внутри `shine-solana/shine/` — читать `shine-solana/shine/AGENTS.md`.
|
||||
- При работе внутри `shine-UI/server-ui/` — читать `shine-UI/AGENTS.md`.
|
||||
- При работе внутри `SHiNE-server/` — читать `SHiNE-server/AGENTS.md`.
|
||||
|
||||
@@ -0,0 +1,47 @@
|
||||
# Правила коммитов, merge и версионирования
|
||||
|
||||
Этот файл является единым источником правил для:
|
||||
- коммитов;
|
||||
- merge в `main`;
|
||||
- изменения версий в `VERSION.properties`;
|
||||
- `git push` из этого репозитория.
|
||||
|
||||
## Язык
|
||||
|
||||
- Пояснения к коммитам, PR и merge-запросам писать на русском языке.
|
||||
|
||||
## Где хранится версия
|
||||
|
||||
- Единый файл версий проекта: `VERSION.properties` в корне репозитория.
|
||||
- Основные поля:
|
||||
- `client.version` — версия клиентского UI.
|
||||
- `server.version` — версия серверной части.
|
||||
|
||||
## Базовое правило для обычных коммитов
|
||||
|
||||
- Перед каждым новым коммитом обязательно обновлять версии в `VERSION.properties`.
|
||||
- Если менялся только UI, увеличивать только `client.version`.
|
||||
- Если менялся только сервер, увеличивать только `server.version`.
|
||||
- Если менялись и UI, и сервер, увеличивать обе версии.
|
||||
- Для обычных коммитов вне `main` использовать стандартный patch-инкремент: `+1` к последнему числовому сегменту.
|
||||
- Пример: `1.2.346` → `1.2.347`.
|
||||
|
||||
## Правило для `main`
|
||||
|
||||
- Ветка `main` предназначена только для стабильных версий.
|
||||
- По умолчанию в `main` нужно не коммитить напрямую, а мержить готовые изменения из рабочей ветки.
|
||||
- Если пользователь просит сделать прямой коммит в `main`, нужно отдельно и явно предупредить, что это обход обычного стабильного процесса, и обязательно переспросить подтверждение.
|
||||
|
||||
## Версионирование при merge или прямом коммите в `main`
|
||||
|
||||
- Для попадания изменений в `main` действует отдельная схема инкремента.
|
||||
- Нужно увеличивать вторую цифру версии и обнулять третью.
|
||||
- Пример: `1.2.346` → `1.3.0`.
|
||||
- Если в наборе изменений менялся только UI, обновлять только `client.version`.
|
||||
- Если в наборе изменений менялся только сервер, обновлять только `server.version`.
|
||||
- Если соответствующая часть не менялась, её версию не трогать.
|
||||
|
||||
## Правила для git commit и git push
|
||||
|
||||
- Обычные коммиты делать стандартным `git commit`; токен для локального коммита не нужен и не используется.
|
||||
- Для операций `git push` при необходимости использовать токен из переменной окружения `$GITEA_TOKEN`.
|
||||
@@ -1,96 +0,0 @@
|
||||
# Production-серверы SHiNE
|
||||
|
||||
## Короткий ответ
|
||||
|
||||
По текущим данным репозитория у SHiNE описаны **два production-контура**:
|
||||
|
||||
- `player@shineup.me`
|
||||
- домен `shineup.me`
|
||||
- IP `178.208.64.62`
|
||||
|
||||
и
|
||||
|
||||
- `player@193.8.215.70`
|
||||
- домен `t.shineup.me`
|
||||
- IP `193.8.215.70`
|
||||
|
||||
Отдельно резервный хост `test.shineup.me` сейчас помечен как временно неработающий.
|
||||
|
||||
## 1. Основной production-хост
|
||||
|
||||
- SSH: `player@shineup.me`
|
||||
- домен: `shineup.me`
|
||||
- IP: `178.208.64.62`
|
||||
- пользователь: `player`
|
||||
- базовый путь: `/home/player`
|
||||
|
||||
Основные каталоги:
|
||||
|
||||
- проект SHiNE: `/home/player/SHiNE`
|
||||
- серверный jar: `/home/player/SHiNE/shine-server/shine-server.jar`
|
||||
- UI: `/home/player/SHiNE/shine-ui`
|
||||
- данные: `/home/player/SHiNE/shine-server/data/`
|
||||
- логи: `/home/player/SHiNE/shine-server/logs/app.log`
|
||||
|
||||
Сервисы:
|
||||
|
||||
- `shine-server.service`
|
||||
- `caddy.service`
|
||||
|
||||
Caddy:
|
||||
|
||||
- активный конфиг: `/etc/caddy/Caddyfile`
|
||||
- UI root: `/home/player/SHiNE/shine-ui`
|
||||
- `/ws` проксируется на `127.0.0.1:7070`
|
||||
|
||||
Deploy:
|
||||
|
||||
- `./gradlew deployServerProduction`
|
||||
- `./gradlew deployUIProduction`
|
||||
|
||||
Правило:
|
||||
|
||||
- любые изменения на `shineup.me` делать только после отдельного подтверждения пользователя.
|
||||
|
||||
## 2. Второй production-сервер
|
||||
|
||||
- SSH: `player@193.8.215.70`
|
||||
- домен: `t.shineup.me`
|
||||
- IP: `193.8.215.70`
|
||||
- пользователь: `player`
|
||||
- базовый путь: `/home/player`
|
||||
|
||||
Роль:
|
||||
|
||||
- второй production-контур SHiNE;
|
||||
- использовать как production-сервер, несмотря на исторические имена deploy-задач
|
||||
`deployServerTest2` / `deployUITest2`.
|
||||
|
||||
Основные каталоги:
|
||||
|
||||
- проект SHiNE: `/home/player/SHiNE`
|
||||
- серверный jar: `/home/player/SHiNE/shine-server/shine-server.jar`
|
||||
- UI: `/home/player/SHiNE/shine-ui`
|
||||
- данные: `/home/player/SHiNE/shine-server/data/`
|
||||
- логи: `/home/player/SHiNE/shine-server/logs/app.log`
|
||||
|
||||
## 3. Связанные публичные production-публикации на том же хосте
|
||||
|
||||
На этом же production-хосте есть отдельная публикация для `shine_payments`:
|
||||
|
||||
- каталог: `/home/player/sites/test-solana-tickets.shineup.me`
|
||||
- домены:
|
||||
- `https://test-solana-tickets.shineup.me`
|
||||
- `https://test-solana-tickets.shiningpeople.ru`
|
||||
|
||||
Это не второй production-хост SHiNE, а отдельный сайт на том же сервере.
|
||||
|
||||
## 4. Какие серверы не считать production
|
||||
|
||||
Не production:
|
||||
|
||||
- `test.shineup.me` (`93.170.12.154`) — резервный хост, временно не работает
|
||||
- `t1.shineup.me`
|
||||
- `t2.shineup.me`
|
||||
- `t3.shineup.me`
|
||||
- `t4.shineup.me`
|
||||
@@ -1,41 +0,0 @@
|
||||
# Deploy
|
||||
|
||||
Подробности о том, где что задеплоено в SHiNE, нужно искать в папке `Deploy/`.
|
||||
|
||||
Эта папка служит краткой картой окружений:
|
||||
|
||||
- [TEST_SERVERS.md](/home/ai/work/SHiNE/SHiNE-server-sha256/Deploy/TEST_SERVERS.md) — тестовые серверы и стенды;
|
||||
- [PRODUCTION_SERVERS.md](/home/ai/work/SHiNE/SHiNE-server-sha256/Deploy/PRODUCTION_SERVERS.md) — production-контур и связанные публичные публикации.
|
||||
|
||||
Ниже краткая сводка.
|
||||
|
||||
## Основные публичные контуры
|
||||
|
||||
- Production SHiNE:
|
||||
- `player@shineup.me`
|
||||
- домен `shineup.me`
|
||||
- IP `178.208.64.62`
|
||||
- Второй production SHiNE:
|
||||
- `player@193.8.215.70`
|
||||
- домен `t.shineup.me`
|
||||
- IP `193.8.215.70`
|
||||
- Временно неработающий резервный test:
|
||||
- `player@93.170.12.154`
|
||||
- домен `test.shineup.me`
|
||||
- IP `93.170.12.154`
|
||||
|
||||
## Отдельный quad-devnet стенд
|
||||
|
||||
На отдельном VPS `178.208.90.249` подняты 4 независимых test/devnet-инстанса:
|
||||
|
||||
- `t1.shineup.me`
|
||||
- `t2.shineup.me`
|
||||
- `t3.shineup.me`
|
||||
- `t4.shineup.me`
|
||||
|
||||
## Важно
|
||||
|
||||
- Production-контура SHiNE сейчас два: `shineup.me` и `t.shineup.me`.
|
||||
- `test.shineup.me` сейчас помечен как временно неработающий резервный test-хост.
|
||||
- `t1..t4.shineup.me` — это отдельные тестовые/devnet-контуры, не production.
|
||||
- Любые изменения на `shineup.me` делать только после отдельного подтверждения пользователя.
|
||||
@@ -1,150 +0,0 @@
|
||||
# Тестовые серверы SHiNE
|
||||
|
||||
Этот файл описывает тестовые контуры, которые сейчас фигурируют в проекте.
|
||||
|
||||
## 1. Исторический `test2`, теперь второй production-сервер
|
||||
|
||||
- SSH: `player@193.8.215.70`
|
||||
- Домен: `t.shineup.me`
|
||||
- IP: `193.8.215.70`
|
||||
- Назначение: второй production-контур SHiNE
|
||||
|
||||
Структура:
|
||||
|
||||
- каталог SHiNE: `/home/player/SHiNE`
|
||||
- сервер: `/home/player/SHiNE/shine-server/shine-server.jar`
|
||||
- UI: `/home/player/SHiNE/shine-ui`
|
||||
- данные: `/home/player/SHiNE/shine-server/data/`
|
||||
- логи: `/home/player/SHiNE/shine-server/logs/app.log`
|
||||
|
||||
Сервисы:
|
||||
|
||||
- `shine-server.service`
|
||||
- `caddy.service`
|
||||
|
||||
Deploy:
|
||||
|
||||
- `./gradlew deployServer`
|
||||
- `./gradlew deployServerTest2`
|
||||
- `./gradlew deployUI`
|
||||
- `./gradlew deployUITest2`
|
||||
|
||||
Примечания:
|
||||
|
||||
- этот хост больше не считать test-контуром;
|
||||
- исторические имена deploy-задач `deployServerTest2` / `deployUITest2` сохранены, но сам хост считать production;
|
||||
- при описании окружений перечислять его как второй production-сервер.
|
||||
|
||||
## 2. Резервный test-сервер, временно неработающий
|
||||
|
||||
- SSH: `player@93.170.12.154`
|
||||
- Домен: `test.shineup.me`
|
||||
- IP: `93.170.12.154`
|
||||
- Назначение: резервный test
|
||||
|
||||
Структура по проектным докам:
|
||||
|
||||
- каталог SHiNE: `/home/player/SHiNE`
|
||||
- сервер: `/home/player/SHiNE/shine-server/shine-server.jar`
|
||||
- UI: `/home/player/SHiNE/shine-ui`
|
||||
- данные: `/home/player/SHiNE/shine-server/data/`
|
||||
- логи: `/home/player/SHiNE/shine-server/logs/app.log`
|
||||
|
||||
Deploy:
|
||||
|
||||
- `./gradlew deployServerTest`
|
||||
- `./gradlew deployUITest`
|
||||
|
||||
Примечания:
|
||||
|
||||
- этот хост резервный;
|
||||
- сейчас помечен как временно неработающий;
|
||||
- использовать его без отдельной причины не нужно.
|
||||
|
||||
## 3. Отдельный quad-devnet стенд `t1..t4`
|
||||
|
||||
- VPS: `178.208.90.249`
|
||||
- пользователь: `player`
|
||||
- назначение: 4 независимых SHiNE-инстанса на Solana `devnet`
|
||||
|
||||
Домены и логины:
|
||||
|
||||
- `server_t1` -> `https://t1.shineup.me`
|
||||
- `server_t2` -> `https://t2.shineup.me`
|
||||
- `server_t3` -> `https://t3.shineup.me`
|
||||
- `server_t4` -> `https://t4.shineup.me`
|
||||
|
||||
Каталоги:
|
||||
|
||||
- `/home/player/t1/server`
|
||||
- `/home/player/t1/UI`
|
||||
- `/home/player/t2/server`
|
||||
- `/home/player/t2/UI`
|
||||
- `/home/player/t3/server`
|
||||
- `/home/player/t3/UI`
|
||||
- `/home/player/t4/server`
|
||||
- `/home/player/t4/UI`
|
||||
|
||||
Подробная памятка на самом VPS:
|
||||
|
||||
- `/home/player/Agents.md`
|
||||
|
||||
Порты и systemd:
|
||||
|
||||
- `t1` -> `7101` -> `shine-t1.service`
|
||||
- `t2` -> `7102` -> `shine-t2.service`
|
||||
- `t3` -> `7103` -> `shine-t3.service`
|
||||
- `t4` -> `7104` -> `shine-t4.service`
|
||||
|
||||
Что важно по конфигу каждого инстанса:
|
||||
|
||||
- отдельный `/home/player/tX/server/application.properties`
|
||||
- `server.port=710X`
|
||||
- `server.SHiNE.login=server_tX`
|
||||
- `db.path=data/shine.sqlite`
|
||||
- `solana.cluster=devnet`
|
||||
- `solana.rpcUrl=https://api.devnet.solana.com`
|
||||
- `server.ui.indexPath=/home/player/tX/UI/index.html`
|
||||
- `server.info.url=https://tX.shineup.me`
|
||||
|
||||
UI каждого инстанса:
|
||||
|
||||
- живёт в отдельной копии `shine-UI`;
|
||||
- использует свой `js/deploy-config.js`;
|
||||
- по умолчанию смотрит именно на свой `tX.shineup.me`.
|
||||
|
||||
Caddy на стенде:
|
||||
|
||||
- конфиг: `/etc/caddy/Caddyfile`
|
||||
- статика: `/home/player/tX/UI`
|
||||
- `/ws` проксируется на `127.0.0.1:710X`
|
||||
|
||||
Operational-нюанс:
|
||||
|
||||
- при одновременных рестартах возможны `HTTP 429` от `api.devnet.solana.com`;
|
||||
- поэтому сервисы `shine-t1..shine-t4` лучше перезапускать по одному, с паузой.
|
||||
|
||||
## 4. Что проверять первым делом
|
||||
|
||||
Для любого test-контура полезны такие быстрые проверки:
|
||||
|
||||
```bash
|
||||
curl -I https://t.shineup.me
|
||||
curl -I https://test.shineup.me
|
||||
curl -I https://t1.shineup.me
|
||||
curl -I https://t2.shineup.me
|
||||
curl -I https://t3.shineup.me
|
||||
curl -I https://t4.shineup.me
|
||||
```
|
||||
|
||||
Для quad-devnet VPS:
|
||||
|
||||
```bash
|
||||
sudo systemctl --no-pager --full status caddy shine-t1 shine-t2 shine-t3 shine-t4
|
||||
```
|
||||
|
||||
Для основного/резервного test:
|
||||
|
||||
```bash
|
||||
sudo systemctl --no-pager --full status shine-server caddy
|
||||
```
|
||||
@@ -1 +0,0 @@
|
||||
|
||||
@@ -1 +0,0 @@
|
||||
|
||||
@@ -1,71 +0,0 @@
|
||||
# Задание для Айдара: навести порядок в инструкциях агентов SHiNE
|
||||
|
||||
## Кратко
|
||||
Нужно согласовать и оформить единый порядок инструкций для Codex/Telegram-агентов в проекте SHiNE, чтобы агенты стабильно понимали структуру проекта, границы ответственности и правила работы с сервером, UI, Solana-модулем, Telegram-ботом и игроками.
|
||||
|
||||
## Зачем это нужно
|
||||
Сейчас проект состоит из нескольких связанных, но разных частей:
|
||||
|
||||
- основной сервер `SHiNE-server/`;
|
||||
- UI `shine-UI/`;
|
||||
- Solana/Anchor-модуль `shine-solana/shine/`;
|
||||
- Telegram-агент-кодер `SHiNE-agent-bot-coder/`;
|
||||
- TURN-сервер;
|
||||
- документация `Dev_Docs/`;
|
||||
- отдельные рабочие папки игроков `Players/`.
|
||||
|
||||
Без явных инструкций агент может путать эти зоны: например, смешать деплой Solana с деплоем сервера, изменить код от имени игрока, не обновить документацию API/DM/блокчейна или неправильно трактовать файл инструкций.
|
||||
|
||||
## Что предлагается сделать
|
||||
1. Утвердить корневой `AGENTS.md` как главный набор правил проекта.
|
||||
2. Проверить и при необходимости уточнить локальный `AGENTS.md` внутри `shine-solana/shine/`.
|
||||
3. Оставить отдельные служебные инструкции Telegram-агента в `SHiNE-agent-bot-coder/AGENT.md`.
|
||||
4. Оставить автоматически читаемые инструкции Telegram-агента в `SHiNE-agent-bot-coder/AGENTS.md`.
|
||||
5. Явно закрепить режим игроков:
|
||||
- игроки могут задавать вопросы, просить анализ, идеи и ТЗ;
|
||||
- игроки не меняют код проекта напрямую;
|
||||
- материалы игроков сохраняются только в `Players/<username>/`.
|
||||
6. Зафиксировать правило: если пользователь говорит «агент MD» или похожую формулировку, считать, что речь про автоматически читаемый `AGENTS.md`.
|
||||
7. Добавить простой процесс согласования изменений инструкций:
|
||||
- Дима или другой участник готовит предложение;
|
||||
- Айдар получает уведомление/заявку;
|
||||
- Айдар отвечает: одобрить, отклонить или попросить доработать;
|
||||
- только после одобрения агент вносит изменения в проектные инструкции.
|
||||
|
||||
## Предлагаемая логика уведомления Айдару
|
||||
Минимальный вариант без сложной разработки:
|
||||
|
||||
1. Агент готовит текст заявки.
|
||||
2. Текст отправляется Айдару в Telegram или в общий рабочий чат.
|
||||
3. В заявке явно указаны варианты ответа:
|
||||
- `одобрить`;
|
||||
- `отклонить`;
|
||||
- `доработать: ...`.
|
||||
4. После ответа Айдара агент либо выполняет согласованные правки, либо фиксирует, что задача отклонена/нужна доработка.
|
||||
|
||||
Более удобный вариант на будущее:
|
||||
|
||||
- добавить в Telegram-бота команду или сценарий согласования задач, например:
|
||||
- `/approve <id>`;
|
||||
- `/reject <id> причина`;
|
||||
- `/revise <id> комментарий`.
|
||||
|
||||
Но для начала достаточно простого текстового согласования через Telegram.
|
||||
|
||||
## Что нужно от Айдара
|
||||
Подтвердить, что такой порядок подходит:
|
||||
|
||||
1. Корневой `AGENTS.md` остается главным правилом проекта.
|
||||
2. Для Solana, Telegram-агента и игроков сохраняются отдельные локальные правила.
|
||||
3. Игроки не меняют код напрямую, а готовят материалы и предложения.
|
||||
4. Изменения инструкций выполняются только после явного одобрения Айдара.
|
||||
5. Уведомления Айдару на первом этапе можно делать простым текстом в Telegram, без отдельной сложной системы заявок.
|
||||
|
||||
## Ожидаемый результат
|
||||
После одобрения:
|
||||
|
||||
- агенты будут стабильнее понимать границы проекта;
|
||||
- снизится риск случайных изменений не в той части системы;
|
||||
- появится понятный порядок согласования задач от игроков;
|
||||
- Айдар будет явно контролировать изменения в инструкциях и правилах работы агентов.
|
||||
|
||||
@@ -1 +0,0 @@
|
||||
|
||||
@@ -1 +0,0 @@
|
||||
|
||||
@@ -1 +0,0 @@
|
||||
|
||||
|
Before Width: | Height: | Size: 31 KiB |
|
Before Width: | Height: | Size: 3.1 KiB |
|
Before Width: | Height: | Size: 105 KiB |
|
Before Width: | Height: | Size: 19 KiB |
|
Before Width: | Height: | Size: 28 KiB |
|
Before Width: | Height: | Size: 93 KiB |
|
Before Width: | Height: | Size: 28 KiB |
|
Before Width: | Height: | Size: 70 KiB |
|
Before Width: | Height: | Size: 28 KiB |
|
Before Width: | Height: | Size: 49 KiB |
|
Before Width: | Height: | Size: 28 KiB |
@@ -1 +0,0 @@
|
||||
|
||||
@@ -1 +0,0 @@
|
||||
|
||||
@@ -1 +0,0 @@
|
||||
|
||||
@@ -1 +0,0 @@
|
||||
|
||||
@@ -1,19 +0,0 @@
|
||||
TELEGRAM_BOT_TOKEN=replace_me
|
||||
OPENAI_API_KEY=replace_me
|
||||
ALLOWED_TELEGRAM_USERNAME=AidarKC
|
||||
ALLOWED_TELEGRAM_PLAYERS=malvviiina:Милана,zodiaktechnika32:Сергей,oidasyda:Иван,blackbyrd1:Ворон,dimasol1:Дима
|
||||
ALLOWED_TELEGRAM_CHANNEL_USERNAME=shine_writing
|
||||
BOT_USERNAME=aidar_su_bot
|
||||
OPENAI_TRANSCRIBE_MODEL=gpt-4o-mini-transcribe
|
||||
TELEGRAM_FILE_DOWNLOAD_TIMEOUT_SECONDS=300
|
||||
OPENAI_TRANSCRIBE_TIMEOUT_SECONDS=900
|
||||
OPENAI_TTS_MODEL=gpt-4o-mini-tts
|
||||
OPENAI_TTS_VOICE=alloy
|
||||
OPENAI_TTS_RESPONSE_FORMAT=opus
|
||||
OPENAI_TTS_TIMEOUT_SECONDS=180
|
||||
OPENAI_TTS_CHUNK_CHARS=3500
|
||||
CODEX_BIN=/home/ai/.cache/JetBrains/IntelliJIdea2026.1/aia/codex/bin/codex-x86_64-unknown-linux-musl
|
||||
CODEX_WORKDIR=/home/ai/work/SHiNE/SHiNE-server-sha256
|
||||
CODEX_TIMEOUT_SECONDS=900
|
||||
MAX_RETRIES=3
|
||||
DATA_DIR=./data
|
||||
@@ -1,5 +0,0 @@
|
||||
.env
|
||||
data/
|
||||
logs/
|
||||
run/
|
||||
__pycache__/
|
||||
@@ -1,81 +0,0 @@
|
||||
# AGENT.md для SHiNE-agent-bot-coder
|
||||
|
||||
Ты запущен как обработчик входящего Telegram-сообщения от пользователя.
|
||||
|
||||
## Контекст
|
||||
- `SHiNE-agent-bot-coder` — локальный Telegram-бот-сервис агента-кодера для работы с этим проектом.
|
||||
- Сервис принимает входящие сообщения от пользователя Telegram, сохраняет историю, ставит задачи в очередь и последовательно запускает Codex CLI в рабочем проекте.
|
||||
- Текстовые сообщения обрабатываются напрямую, voice/audio сначала распознаются через OpenAI transcription, затем передаются как текстовая задача.
|
||||
- История диалога хранится в JSONL-файле, путь передаётся в промпте.
|
||||
- Сообщение может быть текстом или результатом распознавания голосового.
|
||||
- Ответ пойдёт пользователю в Telegram как обычное текстовое сообщение.
|
||||
- Единственная рабочая реализация сервиса — Python-скрипт `py_bot_service.py`; старая Java-реализация удалена как нерабочая и не должна восстанавливаться без отдельного решения Айдара.
|
||||
- В репозитории также есть отдельный Solana/Anchor-модуль `shine-solana/shine/`; он логически связан с SHiNE, но не должен автоматически подключаться к основному серверному deploy без отдельной команды.
|
||||
- Перед изменениями внутри `shine-solana/shine/` читать локальные инструкции `shine-solana/shine/AGENTS.md`; в git не добавлять локальные ключи, `.git`, `.idea`, `.gradle`, `target`, `node_modules`, `test-ledger`, логи, временные run-отчёты и `.env`-конфиги.
|
||||
|
||||
## Авторитет команд и история
|
||||
- Основной пользователь и источник команд — Айдар: `@AidarKC` / `@aidarkc`.
|
||||
- Дополнительно разрешены игроки из whitelist (`ALLOWED_TELEGRAM_PLAYERS`), каждый со своей отдельной историей и рабочей папкой `Players/<username>/`.
|
||||
- Игроки работают в режиме вопросов/анализа/подготовки материалов: в промпте явно задано правило не менять код проекта и писать материалы только в своей папке.
|
||||
- Для неизвестных пользователей в личном чате сервис отвечает вежливым отказом.
|
||||
- В Telegram-канале/группе `@shine_writing` сервис выполняет сообщения только от Айдара, а ответы отправляет в тот же чат.
|
||||
- Если Telegram сообщает о миграции обычной группы в supergroup, сервис должен запомнить новый `chat_id` и отправлять ответы уже туда.
|
||||
- На события подключения/отключения пользователей (join/leave) сервис не отвечает и ничего не отправляет.
|
||||
|
||||
## Очередь и состояние
|
||||
- Входящие задачи записываются в файловую очередь и обрабатываются строго по одной, чтобы не смешивать изменения в проекте.
|
||||
- Сервис ведёт состояние активной задачи и текущего файла истории, а после рестарта продолжает незавершённую обработку с учётом сохранённого состояния.
|
||||
- Истории диалогов хранятся в JSONL по каждому разрешённому username отдельно: `data/history/<username>/`.
|
||||
- Архив истории после `/new`: `data/history/<username>/archive/`.
|
||||
- После `/new` для этого же пользователя должен сбрасываться и контекст продолжения Codex-сессии; следующий запрос запускается как новая сессия, не через resume.
|
||||
- Для просмотра истории игрока открывать файлы в его папке истории по username.
|
||||
- Дедупликация входящих Telegram update нужна, чтобы одно сообщение не попало в обработку повторно.
|
||||
- Если Codex молчит во время активной задачи 2 минуты подряд, сервис отправляет аварийный статус с общим временем работы задачи; при дальнейшем молчании повторяет статус каждые 2 минуты.
|
||||
- После успешной обработки задачи из личного чата Айдара сервис должен отправить публичный итоговый отчёт в группу `@shine_writing`: первым сообщением исходный запрос, вторым сообщением-ответом итоговый ответ Codex. Промежуточные статусы в группу не дублировать.
|
||||
- Для приватных voice/audio-запросов в публичном отчёте первым сообщением отправлять исходный Telegram voice/audio-файл с подписью, где указан распознанный текст. В пользовательском тексте отчёта не показывать Telegram `file_id`.
|
||||
- Озвучивание финальных ответов настраивается персонально для каждого Telegram-пользователя командами `/voice_on`, `/voice_off`; для новых пользователей оно включено по умолчанию.
|
||||
- Адаптация текста перед озвучкой настраивается персонально командами `/voice_rewrite_on`, `/voice_rewrite_off`. Если она включена, сервис перед TTS вызывает дешёвую текстовую модель OpenAI и делает голосовую версию без длинных хэшей, путей, команд и технического шума, сохраняя смысл и порядок исходного ответа.
|
||||
- Режим личных ответов настраивается персонально командами `/single_message_on`, `/single_message_off`: либо одно редактируемое сообщение по этапам, либо отдельные сообщения как раньше.
|
||||
- Команда `/settings` должна сразу показывать текущее состояние всех персональных настроек пользователя и список команд для их изменения.
|
||||
- Если озвучивание включено, после полного текстового финального ответа сервис дополнительно отправляет voice-файл с синтезированной речью через OpenAI TTS даже для текстовых запросов. Voice отправляется в исходный чат, а также в известный личный чат пользователя и в общий чат `@shine_writing`, если они отличаются и доступны. Промежуточные статусы не озвучивать.
|
||||
- Команда `/status` должна показывать состояние очереди и персональные настройки: voice-ответы, адаптацию текста перед озвучкой и режим одного сообщения в личке.
|
||||
|
||||
## Правила голосовой версии ответа
|
||||
- Текстовый финальный ответ должен оставаться полноценным: в нём можно указывать команды, пути, хэши коммитов, номера версий, результаты проверок и другие технические детали.
|
||||
- Голосовую версию финального ответа нужно делать короче и проще для восприятия на слух. Основной механизм — персонально включаемая адаптация текста через дополнительный OpenAI-вызов перед TTS.
|
||||
- В голосовой версии не зачитывать длинные хэши коммитов, токены, file_id, длинные команды, полные пути и другие строки, которые человек всё равно не сможет надёжно запомнить на слух.
|
||||
- Для commit/push в голосовой версии достаточно сказать краткий итог: что коммит сделан, что именно изменено, проверки прошли без ошибок, push выполнен, рабочее дерево чистое.
|
||||
- Если пользователю нужны точные команды, хэши или подробности, они должны оставаться в текстовом ответе.
|
||||
|
||||
## Планы и отложенные фичи
|
||||
- Планы проекта по отложенным фичам хранятся в `TODO/`.
|
||||
- Внутри есть три горизонта:
|
||||
- `near/` - ближайшие планы, обычно сегодня/завтра;
|
||||
- `medium/` - среднесрочные планы, обычно недели или 1-2 месяца;
|
||||
- `far/` - дальнее будущее без понятного срока.
|
||||
- Если пользователь спрашивает, какие есть планы или что можно продолжить, нужно смотреть эти три папки и отвечать кратким списком по горизонтам.
|
||||
- Файлы из `TODO/` не начинать реализовывать без явной команды пользователя.
|
||||
- После реализации фичи, требующей ручной проверки, нужно добавить отдельный файл в `Dev_Docs/Pending_Features/`.
|
||||
|
||||
## Центр задач и предложений
|
||||
- Сервис хранит простые задачи и предложения в `data/task_center/items.json`.
|
||||
- Айдар может смотреть список через `/tasks` или естественные фразы вроде «покажи мои задачи», «покажи задачи Миланы».
|
||||
- Айдар может ставить задачи игрокам фразой вида «поставь задачу Милане: ...».
|
||||
- Игроки могут отправлять предложения Айдару фразой вида `предложение: ...`, `идея: ...` или `заявка: ...`.
|
||||
- Статусы меняются фразами с ID: `одобрить TC-0001`, `отклонить TC-0001`, `доработать TC-0001`, `закрыть TC-0001`.
|
||||
- После финального ответа в личном чате сервис добавляет короткое напоминание, если у пользователя есть активные задачи или предложения.
|
||||
|
||||
## Локальный запуск и systemd
|
||||
- Основной запуск сервиса выполняется Python-скриптом `py_bot_service.py` из папки `SHiNE-agent-bot-coder/`.
|
||||
- Локальные секреты и параметры должны храниться в `.env`, этот файл не коммитится.
|
||||
- Для проверки Codex без Telegram можно использовать self-test режим сервиса.
|
||||
- Для постоянного локального запуска используется user-level systemd service `shine-agent-bot-coder`; скрипты установки лежат в `SHiNE-agent-bot-coder/scripts/systemd/`.
|
||||
- Если меняется логика сервиса, после изменений нужно проверить запуск локально и при необходимости перезапустить user systemd service.
|
||||
- Команда Telegram `/restart` (`/restart_service`) доступна только Айдару и выполняет отложенный рестарт после текущей задачи, до взятия следующей. Аварийный жёсткий рестарт доступен только Айдару командами `/restart_hard`, `/restart_now`, `/restart_force`.
|
||||
|
||||
## Правила ответа
|
||||
- Пиши содержательно и коротко.
|
||||
- Не упоминай внутренние служебные детали, файловую систему и технические логи.
|
||||
- Если запрос требует действий с кодом/проектом, выполняй их в рабочей директории.
|
||||
- Если для ответа данных недостаточно, задай ровно один уточняющий вопрос.
|
||||
- Если была ошибка предыдущего запуска, в промпте будет пометка retry — учти это и продолжи с учётом текущего состояния проекта.
|
||||
@@ -1,25 +0,0 @@
|
||||
# AGENTS
|
||||
|
||||
## Назначение
|
||||
- Это автоматически читаемые инструкции Codex для папки `SHiNE-agent-bot-coder/`.
|
||||
- `SHiNE-agent-bot-coder` — локальный Telegram-бот-сервис агента-кодера для работы с проектом SHiNE.
|
||||
- Если пользователь говорит «агент MD», «агент с MD» или похожим образом про файл инструкций Codex, считать, что имеется в виду `AGENTS.md`.
|
||||
|
||||
## Связанные инструкции
|
||||
- Подробные служебные правила Telegram-обработчика лежат в `AGENT.md`.
|
||||
- `AGENT.md` используется самим сервисом как файл инструкций, который передаётся в промпт обработчика входящих Telegram-сообщений.
|
||||
- При изменении логики сервиса сначала читать `AGENT.md`, затем код `py_bot_service.py`.
|
||||
|
||||
## Планы и задачи
|
||||
- Отложенные задачи проекта лежат в `../TODO/`.
|
||||
- Точка входа по планам: `../TODO/README.md`.
|
||||
- Горизонты планов:
|
||||
- `near/` - ближайшие планы;
|
||||
- `medium/` - среднесрочные планы;
|
||||
- `far/` - дальнее будущее.
|
||||
- Если пользователь спрашивает, какие есть планы или что можно продолжить, кратко перечислять задачи по этим горизонтам.
|
||||
- Не начинать реализацию задач из `TODO` без явной команды пользователя.
|
||||
|
||||
## Проверка после изменений
|
||||
- Если меняется логика Telegram-бота, проверить локальный запуск или self-test, когда это уместно.
|
||||
- Если меняется только документация или инструкции, достаточно проверить, что ссылки на документы актуальны.
|
||||
@@ -1,2 +0,0 @@
|
||||
@AGENTS.md
|
||||
@AGENT.md
|
||||
@@ -1,26 +0,0 @@
|
||||
# Промпты для режима игроков (на согласование)
|
||||
|
||||
## 1) Базовый служебный промпт (добавка к задаче игрока)
|
||||
|
||||
```text
|
||||
Режим игрока (обязательно):
|
||||
- Пользователь: <Имя> (@<username>).
|
||||
- Рабочая папка игрока: <project>/Players/<username>
|
||||
- Код проекта не изменять.
|
||||
- Можно отвечать на вопросы по проекту, предлагать идеи и готовить ТЗ.
|
||||
- Если нужны правки кода, описывать предложение текстом и сохранять материалы только в папке игрока.
|
||||
```
|
||||
|
||||
## 2) Приветственное сообщение игроку (один раз)
|
||||
|
||||
```text
|
||||
Привет, <Имя>.
|
||||
Можно задавать вопросы по проекту, просить анализ, идеи и подготовку готового ТЗ.
|
||||
Команда /new начинает новую сессию и архивирует текущую историю.
|
||||
```
|
||||
|
||||
## 3) Отказ неизвестному пользователю
|
||||
|
||||
```text
|
||||
Извините, доступ к этому агенту пока не выдан. Обратитесь к Айдару.
|
||||
```
|
||||
@@ -1,100 +0,0 @@
|
||||
# SHiNE-agent-bot-coder
|
||||
|
||||
Локальный Telegram-бот-сервис для пользователя `ai`:
|
||||
- принимает сообщения от `@AidarKC`;
|
||||
- поддерживает whitelist игроков (`ALLOWED_TELEGRAM_PLAYERS`) с отдельными историями;
|
||||
- ведёт историю диалога в `JSONL`;
|
||||
- ставит задачи в файловую очередь;
|
||||
- обрабатывает задачи строго последовательно;
|
||||
- поддерживает текстовые и голосовые сообщения (voice/audio через OpenAI transcription);
|
||||
- вызывает Codex CLI и отправляет ответ в Telegram;
|
||||
- в личном чате умеет работать в двух персонально переключаемых режимах: через одно редактируемое статусное сообщение или через отдельные сообщения по этапам;
|
||||
- умеет персонально для каждого пользователя озвучивать финальный ответ через OpenAI TTS;
|
||||
- при рестарте восстанавливает незавершённые задачи;
|
||||
- отправляет аварийный статус только если Codex молчит 2 минуты подряд во время активной задачи;
|
||||
- принимает сообщения из канала/группы `@shine_writing`, выполняет команды только от `@AidarKC`;
|
||||
- учитывает миграцию обычной Telegram-группы в supergroup и перенаправляет ответы на новый `chat_id`.
|
||||
|
||||
Рабочая реализация сервиса — только `py_bot_service.py`. Старая Java-реализация удалена, потому что не заработала и больше не используется.
|
||||
|
||||
## Структура
|
||||
- `.env` — локальные секреты и параметры запуска (не коммитится);
|
||||
- `data/py_queue.jsonl` — очередь Python-сервиса;
|
||||
- `data/py_state.json` — текущее состояние Python-сервиса;
|
||||
- `data/py_processed_updates.log` — дедуп входящих update;
|
||||
- `data/history/<username>/*.jsonl` — активные истории по пользователям;
|
||||
- `data/history/<username>/archive/*.jsonl` — архивы после `/new`.
|
||||
|
||||
## Локальный запуск
|
||||
1. Скопировать пример:
|
||||
- `cp .env.example .env`
|
||||
2. Заполнить секреты в `.env`.
|
||||
- `TELEGRAM_BOT_TOKEN` — токен рабочего Telegram-бота.
|
||||
- `ALLOWED_TELEGRAM_USERNAME` — пользователь, чьи сообщения выполняются как команды.
|
||||
- `ALLOWED_TELEGRAM_PLAYERS` — whitelist игроков в формате `username:Имя,username2:Имя2`.
|
||||
- `ALLOWED_TELEGRAM_CHANNEL_USERNAME` — канал, из которого принимаются `channel_post`; обычные group/supergroup-сообщения обрабатываются как `message`.
|
||||
- `TELEGRAM_API_BASE_URL` — базовый URL Bot API; по умолчанию `https://api.telegram.org`. Для очень больших voice/audio можно поднять локальный `telegram-bot-api` и направить бота туда.
|
||||
- `TELEGRAM_FILE_DOWNLOAD_TIMEOUT_SECONDS` — тайм-аут скачивания voice/audio из Telegram, по умолчанию 300 секунд.
|
||||
- `OPENAI_TRANSCRIBE_TIMEOUT_SECONDS` — тайм-аут распознавания voice/audio в OpenAI, по умолчанию 900 секунд.
|
||||
- `OPENAI_TRANSCRIBE_MAX_UPLOAD_BYTES` — безопасный лимит размера одного куска для OpenAI transcription, по умолчанию `24 MiB`.
|
||||
- `OPENAI_TRANSCRIBE_MAX_CHUNK_SECONDS` — максимальная длина одного куска при длинном аудио, по умолчанию `900` секунд.
|
||||
- `OPENAI_TRANSCRIBE_OVERLAP_SECONDS` — перекрытие соседних кусков для более ровной склейки текста, по умолчанию `2` секунды.
|
||||
- `OPENAI_TRANSCRIBE_REENCODE_BITRATE_KBPS` — битрейт локального пережатия длинного аудио через `ffmpeg`, по умолчанию `24`.
|
||||
- `OPENAI_TRANSCRIBE_FFMPEG_TIMEOUT_SECONDS` — тайм-аут локальной обработки длинного аудио через `ffmpeg`/`ffprobe`, по умолчанию `1800`.
|
||||
- `FFMPEG_BIN` и `FFPROBE_BIN` — пути к локальным бинарям `ffmpeg`/`ffprobe`, если они не лежат в `PATH`.
|
||||
- `OPENAI_TTS_MODEL` — модель синтеза речи, по умолчанию `gpt-4o-mini-tts`.
|
||||
- `OPENAI_TTS_VOICE` — голос синтеза речи, по умолчанию `alloy`.
|
||||
- `OPENAI_TTS_RESPONSE_FORMAT` — аудиоформат для Telegram voice, по умолчанию `opus`.
|
||||
- `OPENAI_TTS_TIMEOUT_SECONDS` — тайм-аут генерации одного фрагмента речи, по умолчанию 180 секунд.
|
||||
- `OPENAI_TTS_CHUNK_CHARS` — максимальный размер одного фрагмента озвучки, по умолчанию 3500 символов.
|
||||
3. Запуск:
|
||||
- `python3 SHiNE-agent-bot-coder/py_bot_service.py`
|
||||
|
||||
## Быстрый self-test Codex (без Telegram)
|
||||
```bash
|
||||
python3 SHiNE-agent-bot-coder/py_bot_service.py --selftest-codex "Ответь одной строкой: Codex работает"
|
||||
```
|
||||
|
||||
## Длинные voice/audio
|
||||
- Если аудио короткое, бот отправляет его в OpenAI как раньше.
|
||||
- Если аудио большое или длинное, бот локально пережимает его через `ffmpeg`, при необходимости режет на куски и распознаёт последовательно.
|
||||
- Если Telegram заранее сообщает большой размер файла, бот больше не отказывается сразу: сначала явно пишет, что пробует скачать файл, затем отдельно сообщает, удалось ли скачивание, и только после успешной загрузки переходит к подготовке аудио и OpenAI.
|
||||
- Для очень больших файлов упираемся не только в OpenAI, но и в лимит обычного облачного Telegram Bot API на скачивание файла ботом. Для таких случаев нужно использовать локальный `telegram-bot-api` сервер и указать его через `TELEGRAM_API_BASE_URL`.
|
||||
|
||||
## Статусы в личке
|
||||
- Для `private`-чата бот поддерживает персональную настройку режима ответа.
|
||||
- По умолчанию он старается не засорять переписку промежуточными сообщениями: создаёт одно статусное сообщение и редактирует его по этапам.
|
||||
- Если включить `/single_message_off`, бот возвращается к старому режиму и отправляет отдельные сообщения по этапам и финальный ответ отдельно.
|
||||
- Если финальный текст в режиме одного сообщения не помещается целиком, бот оставляет первую часть в отредактированном статусном сообщении и отправляет максимум ещё одно дополнительное текстовое сообщение с хвостом ответа.
|
||||
- Голосовой ответ, если он включён, всегда приходит отдельным новым сообщением.
|
||||
|
||||
## Запуск как systemd-сервис
|
||||
Файлы для установки:
|
||||
- `scripts/systemd/shine-agent-bot-coder.service`
|
||||
- `scripts/systemd/install-local-systemd.sh`
|
||||
|
||||
Установка:
|
||||
- `bash SHiNE-agent-bot-coder/scripts/systemd/install-local-systemd.sh`
|
||||
|
||||
Проверка:
|
||||
- `systemctl --user status shine-agent-bot-coder --no-pager`
|
||||
- `journalctl --user -u shine-agent-bot-coder -f`
|
||||
|
||||
Перезапуск после изменений:
|
||||
- `systemctl --user restart shine-agent-bot-coder`
|
||||
|
||||
## Telegram-команды
|
||||
- `/status` — активная задача и размер очереди.
|
||||
- `/settings` — текущие пользовательские настройки и команды для их изменения.
|
||||
- `/queue` — список задач в очереди.
|
||||
- `/stop` — остановить текущую задачу.
|
||||
- `/cancel <id|all>` — удалить задачу по id/префиксу или очистить очередь.
|
||||
- `/new` — архивировать текущую историю, сбросить продолжение Codex-сессии для этого пользователя и начать новый диалог.
|
||||
- `/voice_on` — включить озвучивание финальных ответов для текущего пользователя.
|
||||
- `/voice_off` — выключить озвучивание финальных ответов для текущего пользователя.
|
||||
- `/voice_rewrite_on` — включить адаптацию текста перед озвучкой.
|
||||
- `/voice_rewrite_off` — выключить адаптацию текста перед озвучкой.
|
||||
- `/single_message_on` — вести ответ в личке через одно редактируемое сообщение.
|
||||
- `/single_message_off` — слать отдельные сообщения по этапам и отдельный финальный ответ.
|
||||
- `/restart` или `/restart_service` — отложенный рестарт после текущей задачи, до взятия следующей (только для Айдара).
|
||||
- `/restart_hard` — жёсткий рестарт прямо сейчас (только для Айдара).
|
||||
@@ -1,28 +0,0 @@
|
||||
#!/usr/bin/env bash
|
||||
set -euo pipefail
|
||||
|
||||
ROOT_DIR="/home/ai/work/SHiNE/SHiNE-server-sha256"
|
||||
SERVICE_DIR="${ROOT_DIR}/SHiNE-agent-bot-coder"
|
||||
UNIT_SRC="${SERVICE_DIR}/scripts/systemd/shine-agent-bot-coder.service"
|
||||
UNIT_DST="${HOME}/.config/systemd/user/shine-agent-bot-coder.service"
|
||||
|
||||
echo "[1/6] Проверка python3..."
|
||||
command -v python3 >/dev/null 2>&1 || { echo "python3 не найден"; exit 1; }
|
||||
|
||||
echo "[2/6] Подготовка папки логов..."
|
||||
mkdir -p "${SERVICE_DIR}/logs"
|
||||
|
||||
echo "[3/6] Копирование user systemd unit..."
|
||||
mkdir -p "$(dirname "${UNIT_DST}")"
|
||||
cp "${UNIT_SRC}" "${UNIT_DST}"
|
||||
|
||||
echo "[4/6] daemon-reload..."
|
||||
systemctl --user daemon-reload
|
||||
|
||||
echo "[5/6] enable + start..."
|
||||
systemctl --user enable --now shine-agent-bot-coder
|
||||
|
||||
echo "[6/6] Статус:"
|
||||
systemctl --user status shine-agent-bot-coder --no-pager
|
||||
|
||||
echo "Готово. Логи: journalctl --user -u shine-agent-bot-coder -f"
|
||||
@@ -1,19 +0,0 @@
|
||||
[Unit]
|
||||
Description=SHiNE Agent Bot Coder (Telegram + Codex queue worker)
|
||||
After=network-online.target
|
||||
Wants=network-online.target
|
||||
|
||||
[Service]
|
||||
Type=simple
|
||||
WorkingDirectory=/home/ai/work/SHiNE/SHiNE-server-sha256/SHiNE-agent-bot-coder
|
||||
EnvironmentFile=/home/ai/work/SHiNE/SHiNE-server-sha256/SHiNE-agent-bot-coder/.env
|
||||
ExecStart=/usr/bin/python3 /home/ai/work/SHiNE/SHiNE-server-sha256/SHiNE-agent-bot-coder/py_bot_service.py
|
||||
Restart=always
|
||||
RestartSec=5
|
||||
TimeoutStopSec=20
|
||||
SuccessExitStatus=143 0
|
||||
StandardOutput=append:/home/ai/work/SHiNE/SHiNE-server-sha256/SHiNE-agent-bot-coder/logs/service.log
|
||||
StandardError=append:/home/ai/work/SHiNE/SHiNE-server-sha256/SHiNE-agent-bot-coder/logs/service.log
|
||||
|
||||
[Install]
|
||||
WantedBy=default.target
|
||||
@@ -45,39 +45,20 @@ shine-UI/server-ui.html
|
||||
- `shine_users`: `SHiNEPr1APdAgNBteUyBXcNovaHctpSjUu8oH2ZJdN6`
|
||||
- `shine_payments`: `SHiPmXbM9Fs9khzRUW3TGKsS2W84aqaXTxs3ZkajW9v`
|
||||
|
||||
Подробнее: `Dev_Docs/Инициализация_Solana_регистрации/README.md`
|
||||
Подробнее: `docs/Инициализация_Solana_регистрации/README.md`
|
||||
|
||||
## Синхронизация с партнёрскими серверами
|
||||
|
||||
Сервер должен синхронизировать блоки блокчейна и DM с серверами-партнёрами из `sync_servers`.
|
||||
Детали: `Dev_Docs/Blockchain/sync-between-servers.md`
|
||||
Детали: `docs/Blockchain/sync-between-servers.md`
|
||||
|
||||
## Деплой
|
||||
|
||||
```
|
||||
./gradlew deployServer
|
||||
./gradlew deployUI
|
||||
```
|
||||
|
||||
Default deploy по умолчанию идёт на `t.shineup.me` (`player@193.8.215.70`).
|
||||
|
||||
Production deploy:
|
||||
|
||||
```
|
||||
./gradlew deployServerProduction
|
||||
./gradlew deployUIProduction
|
||||
```
|
||||
|
||||
Любые изменения на `shineup.me` делать только после отдельного явного подтверждения пользователя.
|
||||
|
||||
Резервный test-контур:
|
||||
|
||||
```
|
||||
./gradlew deployServerTest
|
||||
./gradlew deployUITest
|
||||
```
|
||||
|
||||
`test.shineup.me` считается резервным тестовым сервером и в обычный deploy не включается.
|
||||
- Основные инструкции по деплою находятся в `../deploy/AGENTS.md`.
|
||||
- Deploy выполнять shell-скриптами из `../deploy/scripts/`.
|
||||
- Gradle deploy-задачи не использовать: Gradle остаётся для сборки и локального запуска.
|
||||
- Любые изменения на production (`shineup.me`, `server2.shineup.me`) делать только после отдельного явного подтверждения пользователя.
|
||||
- Перед production deploy обязательно обновить/проверить backup в `deploy/backup/archive/`.
|
||||
|
||||
Логи на проде:
|
||||
- `/home/player/SHiNE/shine-server/logs/app.log`
|
||||
|
||||
@@ -59,6 +59,8 @@ public final class AppConfig {
|
||||
public String getParam(String name) {
|
||||
String fromSystem = System.getProperty(name);
|
||||
if (fromSystem != null) return fromSystem;
|
||||
String fromEnv = System.getenv(toEnvName(name));
|
||||
if (fromEnv != null && !fromEnv.isBlank()) return fromEnv.trim();
|
||||
return properties.getProperty(name);
|
||||
}
|
||||
|
||||
@@ -78,4 +80,11 @@ public final class AppConfig {
|
||||
String v = properties.getProperty(name);
|
||||
return v == null ? defaultValue : Boolean.parseBoolean(v);
|
||||
}
|
||||
|
||||
private static String toEnvName(String name) {
|
||||
return name
|
||||
.replace('.', '_')
|
||||
.replace('-', '_')
|
||||
.toUpperCase();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -641,6 +641,7 @@ public final class DatabaseInitializer {
|
||||
origin_session_id TEXT,
|
||||
receipt_ref_base_key TEXT,
|
||||
receipt_ref_type INTEGER,
|
||||
read_at_ms INTEGER,
|
||||
FOREIGN KEY (from_login) REFERENCES solana_users(login),
|
||||
FOREIGN KEY (to_login) REFERENCES solana_users(login)
|
||||
);
|
||||
|
||||
@@ -14,7 +14,7 @@ import java.sql.Statement;
|
||||
public final class SqliteDbController {
|
||||
|
||||
private static volatile SqliteDbController instance;
|
||||
private static final int LATEST_SCHEMA_VERSION = 11;
|
||||
private static final int LATEST_SCHEMA_VERSION = 12;
|
||||
|
||||
private final String jdbcUrl;
|
||||
|
||||
@@ -94,6 +94,7 @@ public final class SqliteDbController {
|
||||
case 9 -> migrateToV9();
|
||||
case 10 -> migrateToV10();
|
||||
case 11 -> migrateToV11();
|
||||
case 12 -> migrateToV12();
|
||||
default -> throw new RuntimeException("Unknown DB migration target version: " + targetVersion);
|
||||
}
|
||||
}
|
||||
@@ -329,6 +330,26 @@ public final class SqliteDbController {
|
||||
}
|
||||
}
|
||||
|
||||
private void migrateToV12() {
|
||||
try (Connection c = DriverManager.getConnection(jdbcUrl);
|
||||
Statement st = c.createStatement()) {
|
||||
c.setAutoCommit(false);
|
||||
try {
|
||||
ensureSignedMessagesReadAtColumn(c, st);
|
||||
backfillSignedMessagesReadAt(st);
|
||||
setSchemaVersion(c, 12);
|
||||
c.commit();
|
||||
} catch (Exception e) {
|
||||
try { c.rollback(); } catch (Exception ignored) {}
|
||||
throw new RuntimeException("DB migration to v12 failed", e);
|
||||
} finally {
|
||||
try { c.setAutoCommit(true); } catch (Exception ignored) {}
|
||||
}
|
||||
} catch (SQLException e) {
|
||||
throw new RuntimeException("DB migration to v12 failed", e);
|
||||
}
|
||||
}
|
||||
|
||||
private static void ensureChat200StateTables(Statement st) throws SQLException {
|
||||
st.executeUpdate("""
|
||||
CREATE TABLE IF NOT EXISTS chat200_state (
|
||||
@@ -463,6 +484,33 @@ public final class SqliteDbController {
|
||||
}
|
||||
}
|
||||
|
||||
private static void ensureSignedMessagesReadAtColumn(Connection c, Statement st) throws SQLException {
|
||||
if (!tableExists(c, "signed_messages_v2")) return;
|
||||
if (!columnExists(c, "signed_messages_v2", "read_at_ms")) {
|
||||
st.executeUpdate("ALTER TABLE signed_messages_v2 ADD COLUMN read_at_ms INTEGER");
|
||||
}
|
||||
}
|
||||
|
||||
private static void backfillSignedMessagesReadAt(Statement st) throws SQLException {
|
||||
st.executeUpdate("""
|
||||
UPDATE signed_messages_v2 AS content
|
||||
SET read_at_ms = (
|
||||
SELECT MIN(receipt.time_ms)
|
||||
FROM signed_messages_v2 AS receipt
|
||||
WHERE receipt.message_type IN (3, 4)
|
||||
AND receipt.receipt_ref_base_key = content.base_key
|
||||
)
|
||||
WHERE content.message_type IN (1, 2)
|
||||
AND (content.read_at_ms IS NULL OR content.read_at_ms <= 0)
|
||||
AND EXISTS (
|
||||
SELECT 1
|
||||
FROM signed_messages_v2 AS receipt
|
||||
WHERE receipt.message_type IN (3, 4)
|
||||
AND receipt.receipt_ref_base_key = content.base_key
|
||||
);
|
||||
""");
|
||||
}
|
||||
|
||||
/**
|
||||
* Временная одноразовая миграция на переходе к SHiNE_DM v1:
|
||||
* старые строки signed_messages_v2 больше не гарантированно совместимы
|
||||
|
||||
@@ -7,10 +7,14 @@ import java.sql.Connection;
|
||||
import java.sql.PreparedStatement;
|
||||
import java.sql.ResultSet;
|
||||
import java.sql.SQLException;
|
||||
import java.sql.Statement;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
|
||||
public final class SignedMessagesV2DAO {
|
||||
private static final int SQLITE_BUSY_MAX_RETRIES = 6;
|
||||
private static final long SQLITE_BUSY_RETRY_BASE_DELAY_MS = 40L;
|
||||
|
||||
public enum ApplyStatus {
|
||||
APPLIED,
|
||||
DUPLICATE_OR_OLDER,
|
||||
@@ -37,6 +41,7 @@ public final class SignedMessagesV2DAO {
|
||||
}
|
||||
|
||||
public ApplyStatus insertIfAbsent(SignedMessageV2Entry e) throws Exception {
|
||||
return withBusyRetry(() -> {
|
||||
try (Connection c = db.getConnection()) {
|
||||
if (isBlockedByConversationDelete(c, e.getFromLogin(), e.getToLogin(), e.getTimeMs())) {
|
||||
return ApplyStatus.BLOCKED_BY_CONVERSATION_TOMBSTONE;
|
||||
@@ -46,17 +51,23 @@ public final class SignedMessagesV2DAO {
|
||||
message_key, base_key, target_login, from_login, to_login,
|
||||
time_ms, nonce, message_type, revision_time_ms, reencrypted_at_ms,
|
||||
raw_block, created_at_ms, source_api, origin_session_id,
|
||||
receipt_ref_base_key, receipt_ref_type
|
||||
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
receipt_ref_base_key, receipt_ref_type, read_at_ms
|
||||
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
""";
|
||||
try (PreparedStatement ps = c.prepareStatement(sql)) {
|
||||
bindSignedMessage(ps, e);
|
||||
return ps.executeUpdate() > 0 ? ApplyStatus.APPLIED : ApplyStatus.DUPLICATE_OR_OLDER;
|
||||
ApplyStatus status = ps.executeUpdate() > 0 ? ApplyStatus.APPLIED : ApplyStatus.DUPLICATE_OR_OLDER;
|
||||
if (status.applied()) {
|
||||
markMessageReadByReceipt(c, e);
|
||||
}
|
||||
return status;
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
public boolean insertPairBothOrNothing(SignedMessageV2Entry first, SignedMessageV2Entry second) throws Exception {
|
||||
return withBusyRetry(() -> {
|
||||
try (Connection c = db.getConnection()) {
|
||||
boolean prevAutoCommit = c.getAutoCommit();
|
||||
c.setAutoCommit(false);
|
||||
@@ -64,6 +75,8 @@ public final class SignedMessagesV2DAO {
|
||||
int insertedFirst = insertStrict(c, first);
|
||||
int insertedSecond = insertStrict(c, second);
|
||||
if (insertedFirst == 1 && insertedSecond == 1) {
|
||||
markMessageReadByReceipt(c, first);
|
||||
markMessageReadByReceipt(c, second);
|
||||
c.commit();
|
||||
return true;
|
||||
}
|
||||
@@ -79,9 +92,11 @@ public final class SignedMessagesV2DAO {
|
||||
c.setAutoCommit(prevAutoCommit);
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
public ApplyStatus upsertContentPair(SignedMessageV2Entry incoming, SignedMessageV2Entry outgoing) throws Exception {
|
||||
return withBusyRetry(() -> {
|
||||
try (Connection c = db.getConnection()) {
|
||||
boolean prevAutoCommit = c.getAutoCommit();
|
||||
c.setAutoCommit(false);
|
||||
@@ -104,6 +119,8 @@ public final class SignedMessagesV2DAO {
|
||||
|
||||
upsertMessage(c, incoming);
|
||||
upsertMessage(c, outgoing);
|
||||
markMessageReadByReceipt(c, incoming);
|
||||
markMessageReadByReceipt(c, outgoing);
|
||||
resetDeliveryRows(c, incoming.getMessageKey());
|
||||
resetDeliveryRows(c, outgoing.getMessageKey());
|
||||
|
||||
@@ -116,9 +133,11 @@ public final class SignedMessagesV2DAO {
|
||||
c.setAutoCommit(prevAutoCommit);
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
public ApplyStatus upsertIncomingCopy(SignedMessageV2Entry incoming) throws Exception {
|
||||
return withBusyRetry(() -> {
|
||||
try (Connection c = db.getConnection()) {
|
||||
boolean prevAutoCommit = c.getAutoCommit();
|
||||
c.setAutoCommit(false);
|
||||
@@ -140,6 +159,7 @@ public final class SignedMessagesV2DAO {
|
||||
}
|
||||
|
||||
upsertMessage(c, incoming);
|
||||
markMessageReadByReceipt(c, incoming);
|
||||
resetDeliveryRows(c, incoming.getMessageKey());
|
||||
c.commit();
|
||||
return ApplyStatus.APPLIED;
|
||||
@@ -150,9 +170,11 @@ public final class SignedMessagesV2DAO {
|
||||
c.setAutoCommit(prevAutoCommit);
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
public ApplyStatus applyDeleteMessage(SignedMessageV2Entry tombstone) throws Exception {
|
||||
return withBusyRetry(() -> {
|
||||
try (Connection c = db.getConnection()) {
|
||||
boolean prevAutoCommit = c.getAutoCommit();
|
||||
c.setAutoCommit(false);
|
||||
@@ -179,9 +201,11 @@ public final class SignedMessagesV2DAO {
|
||||
c.setAutoCommit(prevAutoCommit);
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
public ApplyStatus applyDeleteConversation(SignedMessageV2Entry tombstone) throws Exception {
|
||||
return withBusyRetry(() -> {
|
||||
try (Connection c = db.getConnection()) {
|
||||
boolean prevAutoCommit = c.getAutoCommit();
|
||||
c.setAutoCommit(false);
|
||||
@@ -205,6 +229,7 @@ public final class SignedMessagesV2DAO {
|
||||
c.setAutoCommit(prevAutoCommit);
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
public SignedMessageV2Entry getByMessageKey(String messageKey) throws Exception {
|
||||
@@ -214,7 +239,7 @@ public final class SignedMessagesV2DAO {
|
||||
message_key, base_key, target_login, from_login, to_login,
|
||||
time_ms, nonce, message_type, revision_time_ms, reencrypted_at_ms,
|
||||
raw_block, created_at_ms, source_api, origin_session_id,
|
||||
receipt_ref_base_key, receipt_ref_type
|
||||
receipt_ref_base_key, receipt_ref_type, read_at_ms
|
||||
FROM signed_messages_v2
|
||||
WHERE message_key = ?
|
||||
""";
|
||||
@@ -235,6 +260,12 @@ public final class SignedMessagesV2DAO {
|
||||
}
|
||||
|
||||
public void ensureDeliveryRow(String messageKey, String sessionId, long nowMs) throws Exception {
|
||||
ensureDeliveryRows(messageKey, List.of(sessionId), nowMs);
|
||||
}
|
||||
|
||||
public void ensureDeliveryRows(String messageKey, List<String> sessionIds, long nowMs) throws Exception {
|
||||
if (sessionIds == null || sessionIds.isEmpty()) return;
|
||||
withBusyRetry(() -> {
|
||||
try (Connection c = db.getConnection()) {
|
||||
String sql = """
|
||||
INSERT OR IGNORE INTO signed_message_session_delivery (
|
||||
@@ -242,43 +273,49 @@ public final class SignedMessagesV2DAO {
|
||||
) VALUES (?, ?, 0, NULL, ?)
|
||||
""";
|
||||
try (PreparedStatement ps = c.prepareStatement(sql)) {
|
||||
for (String sessionId : sessionIds) {
|
||||
if (sessionId == null || sessionId.isBlank()) continue;
|
||||
ps.setString(1, messageKey);
|
||||
ps.setString(2, sessionId);
|
||||
ps.setLong(3, nowMs);
|
||||
ps.executeUpdate();
|
||||
ps.addBatch();
|
||||
}
|
||||
ps.executeBatch();
|
||||
}
|
||||
return null;
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
public void markDelivered(String messageKey, String sessionId, long deliveredAtMs) throws Exception {
|
||||
withBusyRetry(() -> {
|
||||
try (Connection c = db.getConnection()) {
|
||||
String insertSql = """
|
||||
INSERT OR IGNORE INTO signed_message_session_delivery (
|
||||
String sql = """
|
||||
INSERT INTO signed_message_session_delivery (
|
||||
message_key, session_id, delivered, delivered_at_ms, created_at_ms
|
||||
) VALUES (?, ?, 0, NULL, ?)
|
||||
) VALUES (?, ?, 1, ?, ?)
|
||||
ON CONFLICT(message_key, session_id) DO UPDATE SET
|
||||
delivered = 1,
|
||||
delivered_at_ms = CASE
|
||||
WHEN signed_message_session_delivery.delivered_at_ms IS NULL THEN excluded.delivered_at_ms
|
||||
WHEN signed_message_session_delivery.delivered_at_ms > excluded.delivered_at_ms THEN excluded.delivered_at_ms
|
||||
ELSE signed_message_session_delivery.delivered_at_ms
|
||||
END
|
||||
""";
|
||||
try (PreparedStatement ps = c.prepareStatement(insertSql)) {
|
||||
try (PreparedStatement ps = c.prepareStatement(sql)) {
|
||||
ps.setString(1, messageKey);
|
||||
ps.setString(2, sessionId);
|
||||
ps.setLong(3, deliveredAtMs);
|
||||
ps.setLong(4, deliveredAtMs);
|
||||
ps.executeUpdate();
|
||||
}
|
||||
|
||||
String updateSql = """
|
||||
UPDATE signed_message_session_delivery
|
||||
SET delivered = 1, delivered_at_ms = ?
|
||||
WHERE message_key = ? AND session_id = ?
|
||||
""";
|
||||
try (PreparedStatement ps = c.prepareStatement(updateSql)) {
|
||||
ps.setLong(1, deliveredAtMs);
|
||||
ps.setString(2, messageKey);
|
||||
ps.setString(3, sessionId);
|
||||
ps.executeUpdate();
|
||||
}
|
||||
return null;
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
public List<SignedMessageV2Entry> listPendingForSession(String login, String sessionId) throws Exception {
|
||||
return withBusyRetry(() -> {
|
||||
try (Connection c = db.getConnection()) {
|
||||
String fillSql = """
|
||||
INSERT OR IGNORE INTO signed_message_session_delivery (
|
||||
@@ -309,7 +346,7 @@ public final class SignedMessagesV2DAO {
|
||||
m.message_key, m.base_key, m.target_login, m.from_login, m.to_login,
|
||||
m.time_ms, m.nonce, m.message_type, m.revision_time_ms, m.reencrypted_at_ms,
|
||||
m.raw_block, m.created_at_ms, m.source_api, m.origin_session_id,
|
||||
m.receipt_ref_base_key, m.receipt_ref_type
|
||||
m.receipt_ref_base_key, m.receipt_ref_type, m.read_at_ms
|
||||
FROM signed_messages_v2 m
|
||||
JOIN signed_message_session_delivery d
|
||||
ON d.message_key = m.message_key
|
||||
@@ -325,6 +362,57 @@ public final class SignedMessagesV2DAO {
|
||||
}
|
||||
return out;
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
public List<SignedMessageV2Entry> listConversationPage(
|
||||
String login,
|
||||
String peerLogin,
|
||||
long beforeTimeMs,
|
||||
String beforeMessageKey,
|
||||
int limit
|
||||
) throws Exception {
|
||||
try (Connection c = db.getConnection()) {
|
||||
String sql = """
|
||||
SELECT
|
||||
message_key, base_key, target_login, from_login, to_login,
|
||||
time_ms, nonce, message_type, revision_time_ms, reencrypted_at_ms,
|
||||
raw_block, created_at_ms, source_api, origin_session_id,
|
||||
receipt_ref_base_key, receipt_ref_type, read_at_ms
|
||||
FROM signed_messages_v2
|
||||
WHERE target_login = ? COLLATE NOCASE
|
||||
AND message_type IN (1, 2)
|
||||
AND (
|
||||
(from_login = ? COLLATE NOCASE AND to_login = ? COLLATE NOCASE)
|
||||
OR (from_login = ? COLLATE NOCASE AND to_login = ? COLLATE NOCASE)
|
||||
)
|
||||
AND (
|
||||
? <= 0
|
||||
OR time_ms < ?
|
||||
OR (time_ms = ? AND (? = '' OR message_key < ?))
|
||||
)
|
||||
ORDER BY time_ms DESC, message_key DESC
|
||||
LIMIT ?
|
||||
""";
|
||||
List<SignedMessageV2Entry> out = new ArrayList<>();
|
||||
try (PreparedStatement ps = c.prepareStatement(sql)) {
|
||||
ps.setString(1, login);
|
||||
ps.setString(2, login);
|
||||
ps.setString(3, peerLogin);
|
||||
ps.setString(4, peerLogin);
|
||||
ps.setString(5, login);
|
||||
ps.setLong(6, beforeTimeMs);
|
||||
ps.setLong(7, beforeTimeMs);
|
||||
ps.setLong(8, beforeTimeMs);
|
||||
ps.setString(9, beforeMessageKey == null ? "" : beforeMessageKey);
|
||||
ps.setString(10, beforeMessageKey == null ? "" : beforeMessageKey);
|
||||
ps.setInt(11, limit);
|
||||
try (ResultSet rs = ps.executeQuery()) {
|
||||
while (rs.next()) out.add(mapRow(rs));
|
||||
}
|
||||
}
|
||||
return out;
|
||||
}
|
||||
}
|
||||
|
||||
private void upsertMessage(Connection c, SignedMessageV2Entry e) throws SQLException {
|
||||
@@ -333,8 +421,8 @@ public final class SignedMessagesV2DAO {
|
||||
message_key, base_key, target_login, from_login, to_login,
|
||||
time_ms, nonce, message_type, revision_time_ms, reencrypted_at_ms,
|
||||
raw_block, created_at_ms, source_api, origin_session_id,
|
||||
receipt_ref_base_key, receipt_ref_type
|
||||
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
receipt_ref_base_key, receipt_ref_type, read_at_ms
|
||||
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
ON CONFLICT(message_key) DO UPDATE SET
|
||||
base_key = excluded.base_key,
|
||||
target_login = excluded.target_login,
|
||||
@@ -350,7 +438,8 @@ public final class SignedMessagesV2DAO {
|
||||
source_api = excluded.source_api,
|
||||
origin_session_id = excluded.origin_session_id,
|
||||
receipt_ref_base_key = excluded.receipt_ref_base_key,
|
||||
receipt_ref_type = excluded.receipt_ref_type
|
||||
receipt_ref_type = excluded.receipt_ref_type,
|
||||
read_at_ms = COALESCE(signed_messages_v2.read_at_ms, excluded.read_at_ms)
|
||||
""";
|
||||
try (PreparedStatement ps = c.prepareStatement(sql)) {
|
||||
bindSignedMessage(ps, e);
|
||||
@@ -358,6 +447,32 @@ public final class SignedMessagesV2DAO {
|
||||
}
|
||||
}
|
||||
|
||||
private void markMessageReadByReceipt(Connection c, SignedMessageV2Entry entry) throws SQLException {
|
||||
if (entry == null) return;
|
||||
int messageType = entry.getMessageType();
|
||||
if (messageType != 3 && messageType != 4) return;
|
||||
String receiptRefBaseKey = String.valueOf(entry.getReceiptRefBaseKey() == null ? "" : entry.getReceiptRefBaseKey()).trim();
|
||||
if (receiptRefBaseKey.isEmpty()) return;
|
||||
long readAtMs = entry.getTimeMs();
|
||||
if (readAtMs <= 0) return;
|
||||
try (PreparedStatement ps = c.prepareStatement("""
|
||||
UPDATE signed_messages_v2
|
||||
SET read_at_ms = CASE
|
||||
WHEN read_at_ms IS NULL OR read_at_ms <= 0 THEN ?
|
||||
WHEN read_at_ms > ? THEN ?
|
||||
ELSE read_at_ms
|
||||
END
|
||||
WHERE base_key = ?
|
||||
AND message_type IN (1, 2)
|
||||
""")) {
|
||||
ps.setLong(1, readAtMs);
|
||||
ps.setLong(2, readAtMs);
|
||||
ps.setLong(3, readAtMs);
|
||||
ps.setString(4, receiptRefBaseKey);
|
||||
ps.executeUpdate();
|
||||
}
|
||||
}
|
||||
|
||||
private RevisionMarker getRevisionMarkerByMessageKey(Connection c, String messageKey) throws SQLException {
|
||||
String sql = """
|
||||
SELECT revision_time_ms, reencrypted_at_ms
|
||||
@@ -536,8 +651,8 @@ public final class SignedMessagesV2DAO {
|
||||
message_key, base_key, target_login, from_login, to_login,
|
||||
time_ms, nonce, message_type, revision_time_ms, reencrypted_at_ms,
|
||||
raw_block, created_at_ms, source_api, origin_session_id,
|
||||
receipt_ref_base_key, receipt_ref_type
|
||||
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
receipt_ref_base_key, receipt_ref_type, read_at_ms
|
||||
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
""";
|
||||
try (PreparedStatement ps = c.prepareStatement(sql)) {
|
||||
bindSignedMessage(ps, e);
|
||||
@@ -563,6 +678,8 @@ public final class SignedMessagesV2DAO {
|
||||
ps.setString(15, e.getReceiptRefBaseKey());
|
||||
if (e.getReceiptRefType() == null) ps.setObject(16, null);
|
||||
else ps.setInt(16, e.getReceiptRefType());
|
||||
if (e.getReadAtMs() == null) ps.setObject(17, null);
|
||||
else ps.setLong(17, e.getReadAtMs());
|
||||
}
|
||||
|
||||
private void bindObjects(PreparedStatement ps, Object... bindValues) throws SQLException {
|
||||
@@ -586,6 +703,45 @@ public final class SignedMessagesV2DAO {
|
||||
return msg.contains("constraint") || msg.contains("unique") || msg.contains("primary key");
|
||||
}
|
||||
|
||||
private boolean isBusyLock(SQLException ex) {
|
||||
Throwable current = ex;
|
||||
while (current != null) {
|
||||
String msg = String.valueOf(current.getMessage()).toLowerCase();
|
||||
if (msg.contains("sqlite_busy") || msg.contains("database is locked") || msg.contains("database table is locked")) {
|
||||
return true;
|
||||
}
|
||||
current = current.getCause();
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
private void sleepBeforeBusyRetry(int attempt) throws SQLException {
|
||||
long delayMs = SQLITE_BUSY_RETRY_BASE_DELAY_MS * (1L << Math.min(attempt, 4));
|
||||
try {
|
||||
Thread.sleep(delayMs);
|
||||
} catch (InterruptedException ie) {
|
||||
Thread.currentThread().interrupt();
|
||||
SQLException sqlEx = new SQLException("Interrupted while retrying SQLite busy lock", ie);
|
||||
throw sqlEx;
|
||||
}
|
||||
}
|
||||
|
||||
private <T> T withBusyRetry(SqlWork<T> work) throws Exception {
|
||||
SQLException lastBusy = null;
|
||||
for (int attempt = 0; attempt < SQLITE_BUSY_MAX_RETRIES; attempt++) {
|
||||
try {
|
||||
return work.run();
|
||||
} catch (SQLException ex) {
|
||||
if (!isBusyLock(ex) || attempt >= SQLITE_BUSY_MAX_RETRIES - 1) {
|
||||
throw ex;
|
||||
}
|
||||
lastBusy = ex;
|
||||
sleepBeforeBusyRetry(attempt);
|
||||
}
|
||||
}
|
||||
throw lastBusy == null ? new SQLException("SQLite busy retry failed") : lastBusy;
|
||||
}
|
||||
|
||||
private int compareMarkers(RevisionMarker left, RevisionMarker right) {
|
||||
int revisionCompare = Long.compare(left.revisionTimeMs, right.revisionTimeMs);
|
||||
if (revisionCompare != 0) return revisionCompare;
|
||||
@@ -611,6 +767,8 @@ public final class SignedMessagesV2DAO {
|
||||
e.setReceiptRefBaseKey(rs.getString("receipt_ref_base_key"));
|
||||
int maybeRefType = rs.getInt("receipt_ref_type");
|
||||
e.setReceiptRefType(rs.wasNull() ? null : maybeRefType);
|
||||
long maybeReadAt = rs.getLong("read_at_ms");
|
||||
e.setReadAtMs(rs.wasNull() ? null : maybeReadAt);
|
||||
return e;
|
||||
}
|
||||
|
||||
@@ -619,4 +777,9 @@ public final class SignedMessagesV2DAO {
|
||||
return new RevisionMarker(entry.getRevisionTimeMs(), entry.getReencryptedAtMs());
|
||||
}
|
||||
}
|
||||
|
||||
@FunctionalInterface
|
||||
private interface SqlWork<T> {
|
||||
T run() throws Exception;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -17,6 +17,7 @@ public class SignedMessageV2Entry {
|
||||
private String originSessionId;
|
||||
private String receiptRefBaseKey;
|
||||
private Integer receiptRefType;
|
||||
private Long readAtMs;
|
||||
|
||||
public String getMessageKey() { return messageKey; }
|
||||
public void setMessageKey(String messageKey) { this.messageKey = messageKey; }
|
||||
@@ -50,4 +51,6 @@ public class SignedMessageV2Entry {
|
||||
public void setReceiptRefBaseKey(String receiptRefBaseKey) { this.receiptRefBaseKey = receiptRefBaseKey; }
|
||||
public Integer getReceiptRefType() { return receiptRefType; }
|
||||
public void setReceiptRefType(Integer receiptRefType) { this.receiptRefType = receiptRefType; }
|
||||
public Long getReadAtMs() { return readAtMs; }
|
||||
public void setReadAtMs(Long readAtMs) { this.readAtMs = readAtMs; }
|
||||
}
|
||||
|
||||
@@ -91,6 +91,7 @@ import server.logic.ws_protocol.JSON.messages.Net_CallInviteBroadcast_Handler;
|
||||
import server.logic.ws_protocol.JSON.messages.Net_CallSignalToSession_Handler;
|
||||
import server.logic.ws_protocol.JSON.messages.Net_DeleteConversation_Handler;
|
||||
import server.logic.ws_protocol.JSON.messages.Net_DeleteMessage_Handler;
|
||||
import server.logic.ws_protocol.JSON.messages.Net_GetDirectMessages_Handler;
|
||||
import server.logic.ws_protocol.JSON.messages.Net_SendSignal_Handler;
|
||||
import server.logic.ws_protocol.JSON.messages.Net_ReceiveIncomingMessage_Handler;
|
||||
import server.logic.ws_protocol.JSON.messages.Net_SendDirectMessage_Handler;
|
||||
@@ -102,6 +103,7 @@ import server.logic.ws_protocol.JSON.messages.entyties.Net_CallInviteBroadcast_R
|
||||
import server.logic.ws_protocol.JSON.messages.entyties.Net_CallSignalToSession_Request;
|
||||
import server.logic.ws_protocol.JSON.messages.entyties.Net_DeleteConversation_Request;
|
||||
import server.logic.ws_protocol.JSON.messages.entyties.Net_DeleteMessage_Request;
|
||||
import server.logic.ws_protocol.JSON.messages.entyties.Net_GetDirectMessages_Request;
|
||||
import server.logic.ws_protocol.JSON.messages.entyties.Net_SendSignal_Request;
|
||||
import server.logic.ws_protocol.JSON.messages.entyties.Net_ReceiveIncomingMessage_Request;
|
||||
import server.logic.ws_protocol.JSON.messages.entyties.Net_SendDirectMessage_Request;
|
||||
@@ -200,6 +202,7 @@ public final class JsonHandlerRegistry {
|
||||
Map.entry("ReceiveIncomingMessage", new Net_ReceiveIncomingMessage_Handler()),
|
||||
Map.entry("DeleteMessage", new Net_DeleteMessage_Handler()),
|
||||
Map.entry("DeleteConversation", new Net_DeleteConversation_Handler()),
|
||||
Map.entry("GetDirectMessages", new Net_GetDirectMessages_Handler()),
|
||||
Map.entry("AckSessionDelivery", new Net_AckSessionDelivery_Handler()),
|
||||
Map.entry("CallInviteBroadcast", new Net_CallInviteBroadcast_Handler()),
|
||||
Map.entry("CallSignalToSession", new Net_CallSignalToSession_Handler()),
|
||||
@@ -280,6 +283,7 @@ public final class JsonHandlerRegistry {
|
||||
Map.entry("ReceiveIncomingMessage", Net_ReceiveIncomingMessage_Request.class),
|
||||
Map.entry("DeleteMessage", Net_DeleteMessage_Request.class),
|
||||
Map.entry("DeleteConversation", Net_DeleteConversation_Request.class),
|
||||
Map.entry("GetDirectMessages", Net_GetDirectMessages_Request.class),
|
||||
Map.entry("AckSessionDelivery", Net_AckSessionDelivery_Request.class),
|
||||
Map.entry("CallInviteBroadcast", Net_CallInviteBroadcast_Request.class),
|
||||
Map.entry("CallSignalToSession", Net_CallSignalToSession_Request.class),
|
||||
|
||||
@@ -10,7 +10,6 @@ import server.logic.ws_protocol.JSON.entyties.Net_Response;
|
||||
import server.logic.ws_protocol.JSON.handlers.JsonMessageHandler;
|
||||
import server.logic.ws_protocol.JSON.handlers.auth.entyties.Net_CreateAuthSession_Request;
|
||||
import server.logic.ws_protocol.JSON.handlers.auth.entyties.Net_CreateAuthSession_Response;
|
||||
import server.logic.ws_protocol.JSON.messages.SignedMessagesRealtime;
|
||||
import server.logic.ws_protocol.JSON.utils.AuthKeyUtils;
|
||||
import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory;
|
||||
import server.logic.ws_protocol.WireCodes;
|
||||
@@ -51,8 +50,6 @@ public class Net_CreateAuthSession__Handler implements JsonMessageHandler {
|
||||
private static final Logger log = LoggerFactory.getLogger(Net_CreateAuthSession__Handler.class);
|
||||
private static final SecureRandom RANDOM = new SecureRandom();
|
||||
private static final long CLOSE_AFTER_ERROR_DELAY_MS = 75L;
|
||||
private static final long SIGNED_DM_BACKLOG_AFTER_AUTH_DELAY_MS = 250L;
|
||||
|
||||
public static final long ALLOWED_SKEW_MS = 30_000L;
|
||||
|
||||
@Override
|
||||
@@ -424,7 +421,6 @@ public class Net_CreateAuthSession__Handler implements JsonMessageHandler {
|
||||
ctx.setAuthenticationStatus(ConnectionContext.AUTH_STATUS_USER);
|
||||
|
||||
ActiveConnectionsRegistry.getInstance().register(ctx);
|
||||
SignedMessagesRealtime.dispatchPendingForSessionAsync(ctx, SIGNED_DM_BACKLOG_AFTER_AUTH_DELAY_MS);
|
||||
|
||||
// --- формируем ответ ---
|
||||
Net_CreateAuthSession_Response resp = new Net_CreateAuthSession_Response();
|
||||
|
||||
@@ -10,7 +10,6 @@ import server.logic.ws_protocol.JSON.entyties.Net_Response;
|
||||
import server.logic.ws_protocol.JSON.handlers.JsonMessageHandler;
|
||||
import server.logic.ws_protocol.JSON.handlers.auth.entyties.Net_SessionLogin_Request;
|
||||
import server.logic.ws_protocol.JSON.handlers.auth.entyties.Net_SessionLogin_Response;
|
||||
import server.logic.ws_protocol.JSON.messages.SignedMessagesRealtime;
|
||||
import server.logic.ws_protocol.JSON.utils.AuthKeyUtils;
|
||||
import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory;
|
||||
import server.logic.ws_protocol.WireCodes;
|
||||
@@ -44,8 +43,6 @@ public class Net_SessionLogin_Handler implements JsonMessageHandler {
|
||||
private static final Logger log = LoggerFactory.getLogger(Net_SessionLogin_Handler.class);
|
||||
|
||||
private static final long ALLOWED_SKEW_MS = 30_000L;
|
||||
private static final long SIGNED_DM_BACKLOG_AFTER_AUTH_DELAY_MS = 250L;
|
||||
|
||||
@Override
|
||||
public Net_Response handle(Net_Request baseReq, ConnectionContext ctx) throws Exception {
|
||||
Net_SessionLogin_Request req = (Net_SessionLogin_Request) baseReq;
|
||||
@@ -302,7 +299,6 @@ public class Net_SessionLogin_Handler implements JsonMessageHandler {
|
||||
ctx.setAuthenticationStatus(ConnectionContext.AUTH_STATUS_USER);
|
||||
|
||||
ActiveConnectionsRegistry.getInstance().register(ctx);
|
||||
SignedMessagesRealtime.dispatchPendingForSessionAsync(ctx, SIGNED_DM_BACKLOG_AFTER_AUTH_DELAY_MS);
|
||||
|
||||
// ответ
|
||||
Net_SessionLogin_Response resp = new Net_SessionLogin_Response();
|
||||
|
||||
@@ -268,7 +268,7 @@ public final class Net_AddBlock_Handler implements JsonMessageHandler {
|
||||
}
|
||||
|
||||
// Репосты временно отключены до будущей реализации.
|
||||
// Точка возврата: Dev_Docs/Future_Features/2026-05-24_1140_репосты_в_каналах_и_тредах.md
|
||||
// Точка возврата: docs/Future_Features/2026-05-24_1140_репосты_в_каналах_и_тредах.md
|
||||
if ((block.type & 0xFFFF) == 1
|
||||
&& (block.subType & 0xFFFF) == (MsgSubType.TEXT_REPOST & 0xFFFF)) {
|
||||
log.warn("AddBlock: repost_disabled (login={}, blockchainName={}, blockNumber={})",
|
||||
|
||||
@@ -0,0 +1,98 @@
|
||||
package server.logic.ws_protocol.JSON.messages;
|
||||
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import server.logic.ws_protocol.JSON.ConnectionContext;
|
||||
import server.logic.ws_protocol.JSON.entyties.Net_Request;
|
||||
import server.logic.ws_protocol.JSON.entyties.Net_Response;
|
||||
import server.logic.ws_protocol.JSON.handlers.JsonMessageHandler;
|
||||
import server.logic.ws_protocol.JSON.messages.entyties.Net_GetDirectMessages_Request;
|
||||
import server.logic.ws_protocol.JSON.messages.entyties.Net_GetDirectMessages_Response;
|
||||
import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory;
|
||||
import server.logic.ws_protocol.WireCodes;
|
||||
import shine.db.dao.SignedMessagesV2DAO;
|
||||
import shine.db.entities.SignedMessageV2Entry;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Base64;
|
||||
import java.util.List;
|
||||
|
||||
public class Net_GetDirectMessages_Handler implements JsonMessageHandler {
|
||||
private static final Logger log = LoggerFactory.getLogger(Net_GetDirectMessages_Handler.class);
|
||||
private static final int DEFAULT_LIMIT = 50;
|
||||
private static final int MAX_LIMIT = 200;
|
||||
|
||||
@Override
|
||||
public Net_Response handle(Net_Request baseRequest, ConnectionContext ctx) {
|
||||
Net_GetDirectMessages_Request req = (Net_GetDirectMessages_Request) baseRequest;
|
||||
if (ctx == null || !ctx.isAuthenticatedUser()) {
|
||||
return NetExceptionResponseFactory.error(req, WireCodes.Status.UNVERIFIED, "NOT_AUTHENTICATED", "Требуется авторизация");
|
||||
}
|
||||
if (req.getPeerLogin() == null || req.getPeerLogin().isBlank()) {
|
||||
return NetExceptionResponseFactory.error(req, WireCodes.Status.BAD_REQUEST, "BAD_FIELDS", "peerLogin обязателен");
|
||||
}
|
||||
|
||||
int limit = req.getLimit() == null ? DEFAULT_LIMIT : req.getLimit();
|
||||
if (limit <= 0 || limit > MAX_LIMIT) {
|
||||
return NetExceptionResponseFactory.error(req, WireCodes.Status.BAD_REQUEST, "BAD_LIMIT", "limit должен быть в диапазоне 1.." + MAX_LIMIT);
|
||||
}
|
||||
|
||||
String login = ctx.getLogin().trim();
|
||||
String peerLogin = req.getPeerLogin().trim();
|
||||
long beforeTimeMs = req.getBeforeTimeMs() == null ? 0L : req.getBeforeTimeMs();
|
||||
String beforeMessageKey = req.getBeforeMessageKey() == null ? "" : req.getBeforeMessageKey().trim();
|
||||
|
||||
try {
|
||||
List<SignedMessageV2Entry> page = SignedMessagesV2DAO.getInstance().listConversationPage(
|
||||
login,
|
||||
peerLogin,
|
||||
beforeTimeMs,
|
||||
beforeMessageKey,
|
||||
limit + 1
|
||||
);
|
||||
|
||||
boolean hasMore = page.size() > limit;
|
||||
if (hasMore) {
|
||||
page = new ArrayList<>(page.subList(0, limit));
|
||||
}
|
||||
|
||||
Net_GetDirectMessages_Response resp = new Net_GetDirectMessages_Response();
|
||||
resp.setOp(req.getOp());
|
||||
resp.setRequestId(req.getRequestId());
|
||||
resp.setStatus(WireCodes.Status.OK);
|
||||
resp.setLogin(login);
|
||||
resp.setPeerLogin(peerLogin);
|
||||
resp.setLimit(limit);
|
||||
resp.setHasMore(hasMore);
|
||||
|
||||
List<Net_GetDirectMessages_Response.MessageItem> items = new ArrayList<>();
|
||||
for (SignedMessageV2Entry entry : page) {
|
||||
Net_GetDirectMessages_Response.MessageItem item = new Net_GetDirectMessages_Response.MessageItem();
|
||||
item.setMessageKey(entry.getMessageKey());
|
||||
item.setBaseKey(entry.getBaseKey());
|
||||
item.setFromLogin(entry.getFromLogin());
|
||||
item.setToLogin(entry.getToLogin());
|
||||
item.setMessageType(entry.getMessageType());
|
||||
item.setTimeMs(entry.getTimeMs());
|
||||
item.setNonce(entry.getNonce());
|
||||
item.setRevisionTimeMs(entry.getRevisionTimeMs());
|
||||
item.setReencryptedAtMs(entry.getReencryptedAtMs());
|
||||
item.setCreatedAtMs(entry.getCreatedAtMs());
|
||||
item.setReadAtMs(entry.getReadAtMs());
|
||||
item.setBlobB64(Base64.getEncoder().encodeToString(entry.getRawBlock()));
|
||||
items.add(item);
|
||||
}
|
||||
resp.setMessages(items);
|
||||
|
||||
if (hasMore && !items.isEmpty()) {
|
||||
Net_GetDirectMessages_Response.MessageItem last = items.get(items.size() - 1);
|
||||
resp.setNextBeforeTimeMs(last.getTimeMs());
|
||||
resp.setNextBeforeMessageKey(last.getMessageKey());
|
||||
}
|
||||
return resp;
|
||||
} catch (Exception e) {
|
||||
log.error("GetDirectMessages failed for login={} peerLogin={}", login, peerLogin, e);
|
||||
return NetExceptionResponseFactory.error(req, WireCodes.Status.INTERNAL_ERROR, "INTERNAL_ERROR", "Внутренняя ошибка сервера");
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -50,12 +50,20 @@ public final class SignedMessagesRealtime {
|
||||
long now = System.currentTimeMillis();
|
||||
for (String targetLogin : targetLoginsForMessage(message)) {
|
||||
List<ActiveSessionEntry> sessions = ActiveSessionsDAO.getInstance().getByLogin(targetLogin);
|
||||
List<String> sessionIdsToTrack = new ArrayList<>();
|
||||
for (ActiveSessionEntry s : sessions) {
|
||||
String sessionId = s.getSessionId();
|
||||
if (excludeSessionId != null && excludeSessionId.equals(sessionId)) {
|
||||
continue;
|
||||
}
|
||||
sessionIdsToTrack.add(sessionId);
|
||||
}
|
||||
SignedMessagesV2DAO.getInstance().ensureDeliveryRows(message.getMessageKey(), sessionIdsToTrack, now);
|
||||
for (ActiveSessionEntry s : sessions) {
|
||||
String sessionId = s.getSessionId();
|
||||
if (excludeSessionId != null && excludeSessionId.equals(sessionId)) {
|
||||
continue;
|
||||
}
|
||||
SignedMessagesV2DAO.getInstance().ensureDeliveryRow(message.getMessageKey(), sessionId, now);
|
||||
boolean deliveredOnline = sendEventToSessionIfOnline(sessionId, targetLogin, message, false);
|
||||
if (deliveredOnline) {
|
||||
counters.wsDelivered++;
|
||||
|
||||
@@ -0,0 +1,19 @@
|
||||
package server.logic.ws_protocol.JSON.messages.entyties;
|
||||
|
||||
import server.logic.ws_protocol.JSON.entyties.Net_Request;
|
||||
|
||||
public class Net_GetDirectMessages_Request extends Net_Request {
|
||||
private String peerLogin;
|
||||
private Integer limit;
|
||||
private Long beforeTimeMs;
|
||||
private String beforeMessageKey;
|
||||
|
||||
public String getPeerLogin() { return peerLogin; }
|
||||
public void setPeerLogin(String peerLogin) { this.peerLogin = peerLogin; }
|
||||
public Integer getLimit() { return limit; }
|
||||
public void setLimit(Integer limit) { this.limit = limit; }
|
||||
public Long getBeforeTimeMs() { return beforeTimeMs; }
|
||||
public void setBeforeTimeMs(Long beforeTimeMs) { this.beforeTimeMs = beforeTimeMs; }
|
||||
public String getBeforeMessageKey() { return beforeMessageKey; }
|
||||
public void setBeforeMessageKey(String beforeMessageKey) { this.beforeMessageKey = beforeMessageKey; }
|
||||
}
|
||||
@@ -0,0 +1,71 @@
|
||||
package server.logic.ws_protocol.JSON.messages.entyties;
|
||||
|
||||
import server.logic.ws_protocol.JSON.entyties.Net_Response;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
|
||||
public class Net_GetDirectMessages_Response extends Net_Response {
|
||||
private String login;
|
||||
private String peerLogin;
|
||||
private int limit;
|
||||
private boolean hasMore;
|
||||
private Long nextBeforeTimeMs;
|
||||
private String nextBeforeMessageKey;
|
||||
private List<MessageItem> messages = new ArrayList<>();
|
||||
|
||||
public String getLogin() { return login; }
|
||||
public void setLogin(String login) { this.login = login; }
|
||||
public String getPeerLogin() { return peerLogin; }
|
||||
public void setPeerLogin(String peerLogin) { this.peerLogin = peerLogin; }
|
||||
public int getLimit() { return limit; }
|
||||
public void setLimit(int limit) { this.limit = limit; }
|
||||
public boolean isHasMore() { return hasMore; }
|
||||
public void setHasMore(boolean hasMore) { this.hasMore = hasMore; }
|
||||
public Long getNextBeforeTimeMs() { return nextBeforeTimeMs; }
|
||||
public void setNextBeforeTimeMs(Long nextBeforeTimeMs) { this.nextBeforeTimeMs = nextBeforeTimeMs; }
|
||||
public String getNextBeforeMessageKey() { return nextBeforeMessageKey; }
|
||||
public void setNextBeforeMessageKey(String nextBeforeMessageKey) { this.nextBeforeMessageKey = nextBeforeMessageKey; }
|
||||
public List<MessageItem> getMessages() { return messages; }
|
||||
public void setMessages(List<MessageItem> messages) { this.messages = messages; }
|
||||
|
||||
public static class MessageItem {
|
||||
private String messageKey;
|
||||
private String baseKey;
|
||||
private String fromLogin;
|
||||
private String toLogin;
|
||||
private int messageType;
|
||||
private long timeMs;
|
||||
private long nonce;
|
||||
private long revisionTimeMs;
|
||||
private long reencryptedAtMs;
|
||||
private long createdAtMs;
|
||||
private Long readAtMs;
|
||||
private String blobB64;
|
||||
|
||||
public String getMessageKey() { return messageKey; }
|
||||
public void setMessageKey(String messageKey) { this.messageKey = messageKey; }
|
||||
public String getBaseKey() { return baseKey; }
|
||||
public void setBaseKey(String baseKey) { this.baseKey = baseKey; }
|
||||
public String getFromLogin() { return fromLogin; }
|
||||
public void setFromLogin(String fromLogin) { this.fromLogin = fromLogin; }
|
||||
public String getToLogin() { return toLogin; }
|
||||
public void setToLogin(String toLogin) { this.toLogin = toLogin; }
|
||||
public int getMessageType() { return messageType; }
|
||||
public void setMessageType(int messageType) { this.messageType = messageType; }
|
||||
public long getTimeMs() { return timeMs; }
|
||||
public void setTimeMs(long timeMs) { this.timeMs = timeMs; }
|
||||
public long getNonce() { return nonce; }
|
||||
public void setNonce(long nonce) { this.nonce = nonce; }
|
||||
public long getRevisionTimeMs() { return revisionTimeMs; }
|
||||
public void setRevisionTimeMs(long revisionTimeMs) { this.revisionTimeMs = revisionTimeMs; }
|
||||
public long getReencryptedAtMs() { return reencryptedAtMs; }
|
||||
public void setReencryptedAtMs(long reencryptedAtMs) { this.reencryptedAtMs = reencryptedAtMs; }
|
||||
public long getCreatedAtMs() { return createdAtMs; }
|
||||
public void setCreatedAtMs(long createdAtMs) { this.createdAtMs = createdAtMs; }
|
||||
public Long getReadAtMs() { return readAtMs; }
|
||||
public void setReadAtMs(Long readAtMs) { this.readAtMs = readAtMs; }
|
||||
public String getBlobB64() { return blobB64; }
|
||||
public void setBlobB64(String blobB64) { this.blobB64 = blobB64; }
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,32 @@
|
||||
plugins {
|
||||
id 'java'
|
||||
}
|
||||
|
||||
group = 'shine'
|
||||
version = '1.0.0'
|
||||
|
||||
java {
|
||||
toolchain {
|
||||
languageVersion = JavaLanguageVersion.of(17)
|
||||
}
|
||||
}
|
||||
|
||||
repositories {
|
||||
mavenCentral()
|
||||
}
|
||||
|
||||
dependencies {
|
||||
implementation project(':shine-server-config')
|
||||
|
||||
implementation 'com.squareup.okhttp3:okhttp:4.12.0'
|
||||
implementation 'com.fasterxml.jackson.core:jackson-databind:2.17.2'
|
||||
implementation 'org.postgresql:postgresql:42.7.7'
|
||||
implementation 'org.bouncycastle:bcprov-jdk18on:1.78.1'
|
||||
implementation 'org.slf4j:slf4j-api:2.0.16'
|
||||
|
||||
testImplementation 'org.junit.jupiter:junit-jupiter:5.11.0'
|
||||
}
|
||||
|
||||
test {
|
||||
useJUnitPlatform()
|
||||
}
|
||||
@@ -0,0 +1,234 @@
|
||||
package sync.config;
|
||||
|
||||
import java.time.Duration;
|
||||
|
||||
public record AppConfig(
|
||||
String rpcUrl,
|
||||
String websocketUrl,
|
||||
String programId,
|
||||
String databaseUrl,
|
||||
String databaseUser,
|
||||
String databasePassword,
|
||||
Duration pollInterval,
|
||||
String commitment
|
||||
) {
|
||||
|
||||
public static final String DEFAULT_PROGRAM_ID =
|
||||
"SHiNEPr1APdAgNBteUyBXcNovaHctpSjUu8oH2ZJdN6";
|
||||
|
||||
public static final String FIXED_COMMITMENT =
|
||||
"confirmed";
|
||||
|
||||
public static final String ENABLED_KEY =
|
||||
"solana.users.sync.enabled";
|
||||
public static final String RPC_URL_KEY =
|
||||
"solana.users.sync.rpcUrl";
|
||||
public static final String WEBSOCKET_URL_KEY =
|
||||
"solana.users.sync.wsUrl";
|
||||
public static final String PROGRAM_ID_KEY =
|
||||
"solana.users.sync.programId";
|
||||
public static final String DATABASE_URL_KEY =
|
||||
"solana.users.sync.databaseUrl";
|
||||
public static final String DATABASE_USER_KEY =
|
||||
"solana.users.sync.dbUser";
|
||||
public static final String DATABASE_PASSWORD_KEY =
|
||||
"solana.users.sync.dbPassword";
|
||||
public static final String POLL_INTERVAL_KEY =
|
||||
"solana.users.sync.pollIntervalSeconds";
|
||||
public static final String LEGACY_SOLANA_RPC_URL_KEY =
|
||||
"solana.rpcUrl";
|
||||
|
||||
public static boolean isEnabled(
|
||||
utils.config.AppConfig serverConfig
|
||||
) {
|
||||
String value =
|
||||
trimToNull(
|
||||
serverConfig.getParam(ENABLED_KEY)
|
||||
);
|
||||
|
||||
if (value == null) {
|
||||
return false;
|
||||
}
|
||||
|
||||
return Boolean.parseBoolean(value);
|
||||
}
|
||||
|
||||
public static AppConfig fromServerConfig(
|
||||
utils.config.AppConfig serverConfig
|
||||
) {
|
||||
|
||||
String rpcUrl =
|
||||
firstRequired(
|
||||
serverConfig,
|
||||
"Solana users sync RPC URL",
|
||||
RPC_URL_KEY,
|
||||
LEGACY_SOLANA_RPC_URL_KEY
|
||||
);
|
||||
|
||||
String websocketUrl =
|
||||
requireParam(
|
||||
serverConfig,
|
||||
WEBSOCKET_URL_KEY
|
||||
);
|
||||
|
||||
String programId =
|
||||
optionalParam(
|
||||
serverConfig,
|
||||
PROGRAM_ID_KEY
|
||||
);
|
||||
|
||||
if (programId == null) {
|
||||
programId =
|
||||
DEFAULT_PROGRAM_ID;
|
||||
}
|
||||
|
||||
String databaseUrl =
|
||||
requireParam(
|
||||
serverConfig,
|
||||
DATABASE_URL_KEY
|
||||
);
|
||||
|
||||
String databaseUser =
|
||||
requireParam(
|
||||
serverConfig,
|
||||
DATABASE_USER_KEY
|
||||
);
|
||||
|
||||
String databasePassword =
|
||||
requireParam(
|
||||
serverConfig,
|
||||
DATABASE_PASSWORD_KEY
|
||||
);
|
||||
|
||||
long pollIntervalSeconds =
|
||||
parsePositiveLong(
|
||||
optionalParam(
|
||||
serverConfig,
|
||||
POLL_INTERVAL_KEY
|
||||
),
|
||||
300L,
|
||||
POLL_INTERVAL_KEY
|
||||
);
|
||||
|
||||
return new AppConfig(
|
||||
rpcUrl,
|
||||
websocketUrl,
|
||||
programId,
|
||||
databaseUrl,
|
||||
databaseUser,
|
||||
databasePassword,
|
||||
Duration.ofSeconds(
|
||||
pollIntervalSeconds
|
||||
),
|
||||
FIXED_COMMITMENT
|
||||
);
|
||||
}
|
||||
|
||||
private static long parsePositiveLong(
|
||||
String rawValue,
|
||||
long defaultValue,
|
||||
String envName
|
||||
) {
|
||||
|
||||
if (rawValue == null) {
|
||||
return defaultValue;
|
||||
}
|
||||
|
||||
try {
|
||||
|
||||
long value =
|
||||
Long.parseLong(
|
||||
rawValue
|
||||
);
|
||||
|
||||
if (value <= 0L) {
|
||||
throw new IllegalArgumentException(
|
||||
envName + " must be > 0"
|
||||
);
|
||||
}
|
||||
|
||||
return value;
|
||||
|
||||
} catch (NumberFormatException exception) {
|
||||
|
||||
throw new IllegalArgumentException(
|
||||
envName + " must be a positive integer",
|
||||
exception
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
private static String firstRequired(
|
||||
utils.config.AppConfig serverConfig,
|
||||
String humanName,
|
||||
String... names
|
||||
) {
|
||||
|
||||
for (String name : names) {
|
||||
String value =
|
||||
optionalParam(
|
||||
serverConfig,
|
||||
name
|
||||
);
|
||||
|
||||
if (value != null) {
|
||||
return value;
|
||||
}
|
||||
}
|
||||
|
||||
throw new IllegalStateException(
|
||||
"Missing required server config: "
|
||||
+ humanName
|
||||
);
|
||||
}
|
||||
|
||||
private static String requireParam(
|
||||
utils.config.AppConfig serverConfig,
|
||||
String name
|
||||
) {
|
||||
|
||||
String value =
|
||||
optionalParam(
|
||||
serverConfig,
|
||||
name
|
||||
);
|
||||
|
||||
if (value == null) {
|
||||
throw new IllegalStateException(
|
||||
"Missing required server config: "
|
||||
+ name
|
||||
);
|
||||
}
|
||||
|
||||
return value;
|
||||
}
|
||||
|
||||
private static String optionalParam(
|
||||
utils.config.AppConfig serverConfig,
|
||||
String name
|
||||
) {
|
||||
|
||||
String value =
|
||||
serverConfig.getParam(name);
|
||||
|
||||
return trimToNull(value);
|
||||
}
|
||||
|
||||
private static String trimToNull(
|
||||
String value
|
||||
) {
|
||||
|
||||
if (value == null) {
|
||||
return null;
|
||||
}
|
||||
|
||||
String trimmed =
|
||||
value.trim();
|
||||
|
||||
if (trimmed.isEmpty()) {
|
||||
return null;
|
||||
}
|
||||
|
||||
return trimmed;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,12 @@
|
||||
package sync.model;
|
||||
|
||||
public record ProgramAccountUpdate(
|
||||
String address,
|
||||
String owner,
|
||||
long lamports,
|
||||
long slot,
|
||||
String dataBase64,
|
||||
boolean executable,
|
||||
Long rentEpoch
|
||||
) {
|
||||
}
|
||||
@@ -0,0 +1,9 @@
|
||||
package sync.model;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
public record SnapshotResult(
|
||||
long snapshotSlot,
|
||||
List<ProgramAccountUpdate> accounts
|
||||
) {
|
||||
}
|
||||
@@ -0,0 +1,822 @@
|
||||
package sync.service;
|
||||
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import sync.codec.ShineUsersCodec;
|
||||
import sync.config.AppConfig;
|
||||
import sync.model.ProgramAccountUpdate;
|
||||
import sync.model.SnapshotResult;
|
||||
import sync.source.AccountUpdateListener;
|
||||
import sync.source.ConnectionListener;
|
||||
import sync.source.rpc.SolanaRpcClient;
|
||||
import sync.source.rpc.SolanaWebSocketClient;
|
||||
import sync.storage.postgres.PostgresStorageRepository;
|
||||
import sync.util.SolanaPdaUtil;
|
||||
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.time.Duration;
|
||||
import java.util.*;
|
||||
import java.util.concurrent.*;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
|
||||
public final class SolanaUsersSyncService
|
||||
implements AutoCloseable {
|
||||
|
||||
private static final int FULL_SNAPSHOT_PROGRESS_STEP =
|
||||
100;
|
||||
|
||||
private static final Logger log =
|
||||
LoggerFactory.getLogger(
|
||||
SolanaUsersSyncService.class
|
||||
);
|
||||
|
||||
private final AppConfig config;
|
||||
private final ObjectMapper mapper;
|
||||
private final PostgresStorageRepository storage;
|
||||
private final SolanaRpcClient rpcClient;
|
||||
private final SolanaWebSocketClient webSocketClient;
|
||||
private final ExecutorService syncExecutor;
|
||||
private final ScheduledExecutorService pollScheduler;
|
||||
private final CompletableFuture<Void> readyFuture =
|
||||
new CompletableFuture<>();
|
||||
private final AtomicBoolean closed =
|
||||
new AtomicBoolean(false);
|
||||
private final AtomicBoolean syncRequested =
|
||||
new AtomicBoolean(false);
|
||||
private final AtomicBoolean syncWorkerScheduled =
|
||||
new AtomicBoolean(false);
|
||||
private final String economyConfigPda;
|
||||
|
||||
private volatile boolean initialSyncCompleted =
|
||||
false;
|
||||
|
||||
public SolanaUsersSyncService(
|
||||
AppConfig config
|
||||
) throws Exception {
|
||||
|
||||
this.config =
|
||||
config;
|
||||
|
||||
this.mapper =
|
||||
new ObjectMapper();
|
||||
|
||||
this.storage =
|
||||
new PostgresStorageRepository(
|
||||
config.databaseUrl(),
|
||||
config.databaseUser(),
|
||||
config.databasePassword(),
|
||||
mapper
|
||||
);
|
||||
|
||||
this.rpcClient =
|
||||
new SolanaRpcClient(
|
||||
config.rpcUrl(),
|
||||
config.programId(),
|
||||
config.commitment()
|
||||
);
|
||||
|
||||
this.webSocketClient =
|
||||
new SolanaWebSocketClient(
|
||||
config.websocketUrl(),
|
||||
config.programId(),
|
||||
config.commitment()
|
||||
);
|
||||
|
||||
this.syncExecutor =
|
||||
Executors.newSingleThreadExecutor(
|
||||
runnable -> {
|
||||
Thread thread =
|
||||
new Thread(
|
||||
runnable,
|
||||
"solana-users-sync-worker"
|
||||
);
|
||||
thread.setDaemon(true);
|
||||
return thread;
|
||||
}
|
||||
);
|
||||
|
||||
this.pollScheduler =
|
||||
Executors.newSingleThreadScheduledExecutor(
|
||||
runnable -> {
|
||||
Thread thread =
|
||||
new Thread(
|
||||
runnable,
|
||||
"solana-users-sync-periodic"
|
||||
);
|
||||
thread.setDaemon(true);
|
||||
return thread;
|
||||
}
|
||||
);
|
||||
|
||||
this.economyConfigPda =
|
||||
SolanaPdaUtil.findProgramAddress(
|
||||
List.of(
|
||||
"shine_users_economy_config"
|
||||
.getBytes(StandardCharsets.UTF_8)
|
||||
),
|
||||
config.programId()
|
||||
);
|
||||
}
|
||||
|
||||
public void start()
|
||||
throws Exception {
|
||||
|
||||
log.info(
|
||||
"Starting sync service. programId={} economyConfigPda={} pollInterval={}",
|
||||
config.programId(),
|
||||
economyConfigPda,
|
||||
config.pollInterval()
|
||||
);
|
||||
|
||||
storage.updateLifecycleState(
|
||||
"STARTING",
|
||||
false,
|
||||
null,
|
||||
null,
|
||||
null
|
||||
);
|
||||
|
||||
webSocketClient.start(
|
||||
new AccountUpdateListener() {
|
||||
@Override
|
||||
public void onAccountUpdate(
|
||||
ProgramAccountUpdate update
|
||||
) {
|
||||
log.debug(
|
||||
"Realtime notification received. address={} slot={}",
|
||||
update.address(),
|
||||
update.slot()
|
||||
);
|
||||
requestSync(
|
||||
"realtime"
|
||||
);
|
||||
}
|
||||
},
|
||||
new ConnectionListener() {
|
||||
@Override
|
||||
public void onConnected(
|
||||
boolean firstConnection
|
||||
) {
|
||||
log.info(
|
||||
"Solana websocket connected. firstConnection={}",
|
||||
firstConnection
|
||||
);
|
||||
requestSync(
|
||||
firstConnection
|
||||
? "initial-connect"
|
||||
: "reconnect"
|
||||
);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onDisconnected(
|
||||
Throwable cause
|
||||
) {
|
||||
log.warn(
|
||||
"Solana websocket disconnected: {}",
|
||||
cause == null
|
||||
? "unknown"
|
||||
: cause.getMessage()
|
||||
);
|
||||
}
|
||||
}
|
||||
);
|
||||
|
||||
long periodSeconds =
|
||||
config.pollInterval()
|
||||
.getSeconds();
|
||||
|
||||
pollScheduler.scheduleWithFixedDelay(
|
||||
() -> requestSync(
|
||||
"periodic"
|
||||
),
|
||||
periodSeconds,
|
||||
periodSeconds,
|
||||
TimeUnit.SECONDS
|
||||
);
|
||||
}
|
||||
|
||||
public void awaitReady()
|
||||
throws Exception {
|
||||
readyFuture.get();
|
||||
}
|
||||
|
||||
public boolean isReady() {
|
||||
return readyFuture.isDone()
|
||||
&& !readyFuture.isCompletedExceptionally();
|
||||
}
|
||||
|
||||
private void requestSync(
|
||||
String reason
|
||||
) {
|
||||
|
||||
if (closed.get()) {
|
||||
return;
|
||||
}
|
||||
|
||||
syncRequested.set(true);
|
||||
|
||||
if (!syncWorkerScheduled.compareAndSet(
|
||||
false,
|
||||
true
|
||||
)) {
|
||||
return;
|
||||
}
|
||||
|
||||
syncExecutor.submit(
|
||||
() -> runSyncLoop(reason)
|
||||
);
|
||||
}
|
||||
|
||||
private void runSyncLoop(
|
||||
String firstReason
|
||||
) {
|
||||
|
||||
String reason =
|
||||
firstReason;
|
||||
|
||||
try {
|
||||
|
||||
while (!closed.get()) {
|
||||
|
||||
boolean shouldRun =
|
||||
syncRequested.getAndSet(false);
|
||||
|
||||
if (!shouldRun) {
|
||||
return;
|
||||
}
|
||||
|
||||
performSync(reason);
|
||||
reason = "coalesced";
|
||||
}
|
||||
|
||||
} catch (Exception exception) {
|
||||
|
||||
log.error(
|
||||
"Sync loop failed",
|
||||
exception
|
||||
);
|
||||
|
||||
try {
|
||||
storage.updateLifecycleState(
|
||||
"FAILED",
|
||||
false,
|
||||
exception.getMessage(),
|
||||
System.currentTimeMillis(),
|
||||
null
|
||||
);
|
||||
} catch (Exception storageException) {
|
||||
log.error(
|
||||
"Failed to persist sync failure state",
|
||||
storageException
|
||||
);
|
||||
}
|
||||
|
||||
readyFuture.completeExceptionally(
|
||||
exception
|
||||
);
|
||||
|
||||
} finally {
|
||||
|
||||
syncWorkerScheduled.set(false);
|
||||
|
||||
if (syncRequested.get()
|
||||
&& !closed.get()
|
||||
&& syncWorkerScheduled.compareAndSet(
|
||||
false,
|
||||
true
|
||||
)) {
|
||||
syncExecutor.submit(
|
||||
() -> runSyncLoop("rescheduled")
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void performSync(
|
||||
String reason
|
||||
) throws Exception {
|
||||
|
||||
long nowMs =
|
||||
System.currentTimeMillis();
|
||||
|
||||
PostgresStorageRepository.SyncStateSnapshot state =
|
||||
storage.loadState();
|
||||
|
||||
log.info(
|
||||
"Starting history sync. reason={} lastSeenSignature={}",
|
||||
reason,
|
||||
state.lastSeenSignature()
|
||||
);
|
||||
|
||||
storage.updateLifecycleState(
|
||||
initialSyncCompleted
|
||||
? "SYNCING"
|
||||
: "BOOTSTRAPPING",
|
||||
false,
|
||||
null,
|
||||
nowMs,
|
||||
null
|
||||
);
|
||||
|
||||
SolanaRpcClient.SignatureFetchResult fetchResult =
|
||||
rpcClient.getSignaturesForAddressSince(
|
||||
economyConfigPda,
|
||||
state.lastSeenSignature()
|
||||
);
|
||||
|
||||
if (state.lastSeenSignature() != null
|
||||
&& !fetchResult.anchorFound()) {
|
||||
|
||||
log.error(
|
||||
"History anchor signature not found anymore: {}. Running current-state full snapshot fallback.",
|
||||
state.lastSeenSignature()
|
||||
);
|
||||
|
||||
runFullSnapshotFallback(
|
||||
state,
|
||||
fetchResult,
|
||||
nowMs
|
||||
);
|
||||
|
||||
markReadyAfterSync(
|
||||
state,
|
||||
nowMs
|
||||
);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
if (fetchResult.signatures().isEmpty()) {
|
||||
|
||||
PostgresStorageRepository.SyncStateSnapshot newState =
|
||||
new PostgresStorageRepository.SyncStateSnapshot(
|
||||
"READY",
|
||||
true,
|
||||
nowMs,
|
||||
nowMs,
|
||||
state.lastSeenSignature(),
|
||||
state.lastSeenSlot(),
|
||||
state.lastRelevantSignature(),
|
||||
state.lastRelevantSlot(),
|
||||
null,
|
||||
state.economyConfigState(),
|
||||
nowMs
|
||||
);
|
||||
|
||||
storage.applyHistoryBatch(
|
||||
List.of(),
|
||||
List.of(),
|
||||
newState
|
||||
);
|
||||
|
||||
log.info(
|
||||
"History sync completed with no new transactions."
|
||||
);
|
||||
|
||||
markReadyAfterSync(
|
||||
newState,
|
||||
nowMs
|
||||
);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
List<SolanaRpcClient.SignatureRecord> chronologicalSignatures =
|
||||
new ArrayList<>(
|
||||
fetchResult.signatures()
|
||||
);
|
||||
|
||||
Collections.reverse(
|
||||
chronologicalSignatures
|
||||
);
|
||||
|
||||
List<ParsedTxEnvelope> envelopes =
|
||||
new ArrayList<>();
|
||||
|
||||
Set<String> updatePdaAddresses =
|
||||
new LinkedHashSet<>();
|
||||
|
||||
for (SolanaRpcClient.SignatureRecord signatureRecord : chronologicalSignatures) {
|
||||
|
||||
JsonNode transaction =
|
||||
rpcClient.getTransactionJsonParsed(
|
||||
signatureRecord.signature()
|
||||
);
|
||||
|
||||
ParsedTxEnvelope envelope =
|
||||
parseTransactionEnvelope(
|
||||
signatureRecord,
|
||||
transaction
|
||||
);
|
||||
|
||||
envelopes.add(
|
||||
envelope
|
||||
);
|
||||
|
||||
if (envelope.parsedInstruction() != null
|
||||
&& envelope.parsedInstruction().kind() == ShineUsersCodec.TxKind.UPDATE_USER_PDA
|
||||
&& envelope.parsedInstruction().affectedPdaAddress() != null) {
|
||||
updatePdaAddresses.add(
|
||||
envelope.parsedInstruction()
|
||||
.affectedPdaAddress()
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Map<String, ShineUsersCodec.UserPdaSnapshot> currentSnapshots =
|
||||
storage.getCurrentSnapshots(
|
||||
updatePdaAddresses
|
||||
);
|
||||
|
||||
List<PostgresStorageRepository.TxHistoryEntry> txEntries =
|
||||
new ArrayList<>();
|
||||
|
||||
List<ShineUsersCodec.UserPdaSnapshot> snapshotsToPersist =
|
||||
new ArrayList<>();
|
||||
|
||||
ShineUsersCodec.EconomyConfigState economyState =
|
||||
state.economyConfigState();
|
||||
|
||||
String lastRelevantSignature =
|
||||
state.lastRelevantSignature();
|
||||
|
||||
Long lastRelevantSlot =
|
||||
state.lastRelevantSlot();
|
||||
|
||||
for (ParsedTxEnvelope envelope : envelopes) {
|
||||
|
||||
ShineUsersCodec.ParsedInstruction parsedInstruction =
|
||||
envelope.parsedInstruction();
|
||||
|
||||
String txKind =
|
||||
parsedInstruction == null
|
||||
? "failed_or_unavailable"
|
||||
: parsedInstruction.kind().name();
|
||||
|
||||
boolean relevant =
|
||||
false;
|
||||
|
||||
String affectedPdaAddress =
|
||||
null;
|
||||
|
||||
String affectedLogin =
|
||||
null;
|
||||
|
||||
if (parsedInstruction != null) {
|
||||
|
||||
if (parsedInstruction.kind() == ShineUsersCodec.TxKind.INIT_USERS_ECONOMY_CONFIG
|
||||
|| parsedInstruction.kind() == ShineUsersCodec.TxKind.UPDATE_USERS_ECONOMY_CONFIG) {
|
||||
economyState =
|
||||
parsedInstruction.economyConfigState();
|
||||
}
|
||||
|
||||
if (parsedInstruction.relevant()
|
||||
&& parsedInstruction.userPdaMutation() != null) {
|
||||
|
||||
relevant = true;
|
||||
affectedPdaAddress = parsedInstruction.affectedPdaAddress();
|
||||
affectedLogin = parsedInstruction.affectedLogin();
|
||||
|
||||
ShineUsersCodec.UserPdaSnapshot snapshot;
|
||||
|
||||
if (parsedInstruction.kind() == ShineUsersCodec.TxKind.CREATE_USER_PDA) {
|
||||
|
||||
if (economyState == null) {
|
||||
economyState =
|
||||
ShineUsersCodec.EconomyConfigState.initial();
|
||||
log.warn(
|
||||
"Economy config state was absent while processing create tx {}. Falling back to initial constants.",
|
||||
envelope.signatureRecord().signature()
|
||||
);
|
||||
}
|
||||
|
||||
snapshot =
|
||||
ShineUsersCodec.buildCreateSnapshot(
|
||||
parsedInstruction.userPdaMutation(),
|
||||
economyState,
|
||||
envelope.signatureRecord().signature(),
|
||||
envelope.signatureRecord().slot()
|
||||
);
|
||||
|
||||
} else {
|
||||
|
||||
ShineUsersCodec.UserPdaSnapshot previous =
|
||||
currentSnapshots.get(
|
||||
affectedPdaAddress
|
||||
);
|
||||
|
||||
if (previous == null) {
|
||||
throw new IllegalStateException(
|
||||
"Missing previous snapshot for update PDA " +
|
||||
affectedPdaAddress
|
||||
);
|
||||
}
|
||||
|
||||
snapshot =
|
||||
ShineUsersCodec.buildUpdateSnapshot(
|
||||
parsedInstruction.userPdaMutation(),
|
||||
previous,
|
||||
envelope.signatureRecord().signature(),
|
||||
envelope.signatureRecord().slot()
|
||||
);
|
||||
}
|
||||
|
||||
currentSnapshots.put(
|
||||
snapshot.pdaAddress(),
|
||||
snapshot
|
||||
);
|
||||
|
||||
snapshotsToPersist.add(
|
||||
snapshot
|
||||
);
|
||||
|
||||
lastRelevantSignature =
|
||||
envelope.signatureRecord().signature();
|
||||
|
||||
lastRelevantSlot =
|
||||
envelope.signatureRecord().slot();
|
||||
}
|
||||
}
|
||||
|
||||
txEntries.add(
|
||||
new PostgresStorageRepository.TxHistoryEntry(
|
||||
envelope.signatureRecord().signature(),
|
||||
envelope.signatureRecord().slot(),
|
||||
envelope.signatureRecord().blockTime(),
|
||||
txKind,
|
||||
relevant,
|
||||
affectedPdaAddress,
|
||||
affectedLogin,
|
||||
envelope.rawTransactionJson(),
|
||||
nowMs
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
SolanaRpcClient.SignatureRecord newestSeen =
|
||||
fetchResult.signatures()
|
||||
.get(0);
|
||||
|
||||
PostgresStorageRepository.SyncStateSnapshot newState =
|
||||
new PostgresStorageRepository.SyncStateSnapshot(
|
||||
"READY",
|
||||
true,
|
||||
nowMs,
|
||||
nowMs,
|
||||
newestSeen.signature(),
|
||||
newestSeen.slot(),
|
||||
lastRelevantSignature,
|
||||
lastRelevantSlot,
|
||||
null,
|
||||
economyState,
|
||||
nowMs
|
||||
);
|
||||
|
||||
storage.applyHistoryBatch(
|
||||
txEntries,
|
||||
snapshotsToPersist,
|
||||
newState
|
||||
);
|
||||
|
||||
log.info(
|
||||
"History sync completed. txCount={} relevantCount={} latestSignature={}",
|
||||
txEntries.size(),
|
||||
snapshotsToPersist.size(),
|
||||
newestSeen.signature()
|
||||
);
|
||||
|
||||
markReadyAfterSync(
|
||||
newState,
|
||||
nowMs
|
||||
);
|
||||
}
|
||||
|
||||
private void runFullSnapshotFallback(
|
||||
PostgresStorageRepository.SyncStateSnapshot state,
|
||||
SolanaRpcClient.SignatureFetchResult fetchResult,
|
||||
long nowMs
|
||||
) throws Exception {
|
||||
|
||||
log.warn(
|
||||
"Starting full snapshot fallback because incremental history anchor is unavailable."
|
||||
);
|
||||
|
||||
SnapshotResult snapshotResult =
|
||||
rpcClient.loadFullSnapshot();
|
||||
|
||||
log.info(
|
||||
"Full snapshot downloaded. snapshotSlot={} rawAccounts={}",
|
||||
snapshotResult.snapshotSlot(),
|
||||
snapshotResult.accounts().size()
|
||||
);
|
||||
|
||||
List<ShineUsersCodec.UserPdaSnapshot> currentSnapshots =
|
||||
new ArrayList<>();
|
||||
|
||||
int processedAccounts =
|
||||
0;
|
||||
|
||||
for (ProgramAccountUpdate account : snapshotResult.accounts()) {
|
||||
|
||||
processedAccounts++;
|
||||
|
||||
try {
|
||||
currentSnapshots.add(
|
||||
ShineUsersCodec.parseUserPdaAccount(
|
||||
account.address(),
|
||||
account.slot(),
|
||||
account.dataBase64(),
|
||||
state.lastSeenSignature()
|
||||
)
|
||||
);
|
||||
} catch (Exception ignored) {
|
||||
}
|
||||
|
||||
if (processedAccounts == 1
|
||||
|| processedAccounts % FULL_SNAPSHOT_PROGRESS_STEP == 0
|
||||
|| processedAccounts == snapshotResult.accounts().size()) {
|
||||
log.info(
|
||||
"Full snapshot parse progress: {}/{} accounts, {} user PDA snapshots accepted.",
|
||||
processedAccounts,
|
||||
snapshotResult.accounts().size(),
|
||||
currentSnapshots.size()
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
String newestSignature =
|
||||
fetchResult.signatures().isEmpty()
|
||||
? state.lastSeenSignature()
|
||||
: fetchResult.signatures().get(0).signature();
|
||||
|
||||
Long newestSlot =
|
||||
fetchResult.signatures().isEmpty()
|
||||
? state.lastSeenSlot()
|
||||
: fetchResult.signatures().get(0).slot();
|
||||
|
||||
PostgresStorageRepository.SyncStateSnapshot newState =
|
||||
new PostgresStorageRepository.SyncStateSnapshot(
|
||||
"READY",
|
||||
true,
|
||||
nowMs,
|
||||
nowMs,
|
||||
newestSignature,
|
||||
newestSlot,
|
||||
state.lastRelevantSignature(),
|
||||
state.lastRelevantSlot(),
|
||||
"history_anchor_missing_full_snapshot_fallback",
|
||||
state.economyConfigState(),
|
||||
nowMs
|
||||
);
|
||||
|
||||
storage.replaceCurrentFromFullSnapshot(
|
||||
currentSnapshots,
|
||||
newState
|
||||
);
|
||||
|
||||
log.warn(
|
||||
"Full snapshot fallback completed. currentSnapshots={} newestSignature={} newestSlot={}",
|
||||
currentSnapshots.size(),
|
||||
newestSignature,
|
||||
newestSlot
|
||||
);
|
||||
}
|
||||
|
||||
private void markReadyAfterSync(
|
||||
PostgresStorageRepository.SyncStateSnapshot state,
|
||||
long nowMs
|
||||
) throws Exception {
|
||||
|
||||
initialSyncCompleted = true;
|
||||
|
||||
if (!readyFuture.isDone()) {
|
||||
readyFuture.complete(null);
|
||||
log.info(
|
||||
"Sync service entered READY state."
|
||||
);
|
||||
}
|
||||
|
||||
storage.updateLifecycleState(
|
||||
"READY",
|
||||
true,
|
||||
state.lastError(),
|
||||
nowMs,
|
||||
nowMs
|
||||
);
|
||||
}
|
||||
|
||||
private ParsedTxEnvelope parseTransactionEnvelope(
|
||||
SolanaRpcClient.SignatureRecord signatureRecord,
|
||||
JsonNode transaction
|
||||
) {
|
||||
|
||||
if (signatureRecord.failed()) {
|
||||
return new ParsedTxEnvelope(
|
||||
signatureRecord,
|
||||
new ShineUsersCodec.ParsedInstruction(
|
||||
ShineUsersCodec.TxKind.OTHER,
|
||||
false,
|
||||
null,
|
||||
null,
|
||||
null,
|
||||
null
|
||||
),
|
||||
transaction == null
|
||||
? "{\"failed\":true,\"error\":" +
|
||||
String.valueOf(signatureRecord.errorJson()) + "}"
|
||||
: transaction.toString()
|
||||
);
|
||||
}
|
||||
|
||||
if (transaction == null
|
||||
|| transaction.isNull()) {
|
||||
return new ParsedTxEnvelope(
|
||||
signatureRecord,
|
||||
null,
|
||||
"{\"transaction\":null}"
|
||||
);
|
||||
}
|
||||
|
||||
JsonNode instructions =
|
||||
transaction.path("transaction")
|
||||
.path("message")
|
||||
.path("instructions");
|
||||
|
||||
if (instructions.isArray()) {
|
||||
for (JsonNode instruction : instructions) {
|
||||
ShineUsersCodec.ParsedInstruction parsedInstruction =
|
||||
ShineUsersCodec.parseShineUsersInstruction(
|
||||
instruction,
|
||||
config.programId()
|
||||
);
|
||||
|
||||
if (parsedInstruction != null) {
|
||||
return new ParsedTxEnvelope(
|
||||
signatureRecord,
|
||||
parsedInstruction,
|
||||
transaction.toString()
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return new ParsedTxEnvelope(
|
||||
signatureRecord,
|
||||
new ShineUsersCodec.ParsedInstruction(
|
||||
ShineUsersCodec.TxKind.OTHER,
|
||||
false,
|
||||
null,
|
||||
null,
|
||||
null,
|
||||
null
|
||||
),
|
||||
transaction.toString()
|
||||
);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
|
||||
if (!closed.compareAndSet(
|
||||
false,
|
||||
true
|
||||
)) {
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
pollScheduler.shutdownNow();
|
||||
} catch (Exception ignored) {
|
||||
}
|
||||
|
||||
try {
|
||||
webSocketClient.close();
|
||||
} catch (Exception ignored) {
|
||||
}
|
||||
|
||||
try {
|
||||
syncExecutor.shutdownNow();
|
||||
} catch (Exception ignored) {
|
||||
}
|
||||
|
||||
try {
|
||||
rpcClient.close();
|
||||
} catch (Exception ignored) {
|
||||
}
|
||||
|
||||
try {
|
||||
storage.close();
|
||||
} catch (Exception ignored) {
|
||||
}
|
||||
}
|
||||
|
||||
private record ParsedTxEnvelope(
|
||||
SolanaRpcClient.SignatureRecord signatureRecord,
|
||||
ShineUsersCodec.ParsedInstruction parsedInstruction,
|
||||
String rawTransactionJson
|
||||
) {
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
package sync.source;
|
||||
|
||||
import sync.model.ProgramAccountUpdate;
|
||||
|
||||
@FunctionalInterface
|
||||
public interface AccountUpdateListener {
|
||||
|
||||
void onAccountUpdate(
|
||||
ProgramAccountUpdate update
|
||||
);
|
||||
}
|
||||
@@ -0,0 +1,12 @@
|
||||
package sync.source;
|
||||
|
||||
public interface ConnectionListener {
|
||||
|
||||
void onConnected(
|
||||
boolean firstConnection
|
||||
);
|
||||
|
||||
void onDisconnected(
|
||||
Throwable cause
|
||||
);
|
||||
}
|
||||
@@ -0,0 +1,737 @@
|
||||
package sync.source.rpc;
|
||||
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import okhttp3.*;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import sync.model.ProgramAccountUpdate;
|
||||
import sync.model.SnapshotResult;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.time.Duration;
|
||||
import java.util.*;
|
||||
|
||||
public final class SolanaRpcClient
|
||||
implements AutoCloseable {
|
||||
|
||||
private static final Logger log =
|
||||
LoggerFactory.getLogger(
|
||||
SolanaRpcClient.class
|
||||
);
|
||||
|
||||
private static final MediaType JSON =
|
||||
MediaType.get("application/json");
|
||||
|
||||
private static final int SIGNATURE_PAGE_SIZE =
|
||||
1000;
|
||||
|
||||
private static final int ACCOUNT_BATCH_SIZE =
|
||||
100;
|
||||
|
||||
private final String rpcUrl;
|
||||
private final String programId;
|
||||
private final String commitment;
|
||||
private final ObjectMapper mapper =
|
||||
new ObjectMapper();
|
||||
private final OkHttpClient httpClient =
|
||||
new OkHttpClient.Builder()
|
||||
.callTimeout(Duration.ofMinutes(2))
|
||||
.build();
|
||||
|
||||
public SolanaRpcClient(
|
||||
String rpcUrl,
|
||||
String programId,
|
||||
String commitment
|
||||
) {
|
||||
this.rpcUrl = rpcUrl;
|
||||
this.programId = programId;
|
||||
this.commitment = commitment;
|
||||
}
|
||||
|
||||
public SnapshotResult loadFullSnapshot()
|
||||
throws IOException {
|
||||
|
||||
log.info(
|
||||
"Requesting full snapshot via getProgramAccounts. programId={}",
|
||||
programId
|
||||
);
|
||||
|
||||
Map<String, Object> payload =
|
||||
Map.of(
|
||||
"jsonrpc", "2.0",
|
||||
"id", 100,
|
||||
"method", "getProgramAccounts",
|
||||
"params", List.of(
|
||||
programId,
|
||||
Map.of(
|
||||
"encoding", "base64",
|
||||
"commitment", commitment,
|
||||
"withContext", true
|
||||
)
|
||||
)
|
||||
);
|
||||
|
||||
JsonNode root =
|
||||
executeRpc(payload);
|
||||
|
||||
log.info(
|
||||
"Full snapshot RPC response received. Parsing account list..."
|
||||
);
|
||||
|
||||
JsonNode result =
|
||||
root.path("result");
|
||||
|
||||
long snapshotSlot =
|
||||
result.path("context")
|
||||
.path("slot")
|
||||
.asLong(-1);
|
||||
|
||||
if (snapshotSlot < 0) {
|
||||
throw new IOException(
|
||||
"Missing context.slot in getProgramAccounts"
|
||||
);
|
||||
}
|
||||
|
||||
JsonNode values =
|
||||
result.path("value");
|
||||
|
||||
if (!values.isArray()) {
|
||||
throw new IOException(
|
||||
"Unexpected getProgramAccounts response"
|
||||
);
|
||||
}
|
||||
|
||||
List<ProgramAccountUpdate> accounts =
|
||||
new ArrayList<>();
|
||||
|
||||
for (JsonNode item : values) {
|
||||
|
||||
JsonNode account =
|
||||
item.path("account");
|
||||
|
||||
JsonNode data =
|
||||
account.path("data");
|
||||
|
||||
if (!data.isArray()
|
||||
|| data.isEmpty()) {
|
||||
continue;
|
||||
}
|
||||
|
||||
accounts.add(
|
||||
parseAccount(
|
||||
item.path("pubkey")
|
||||
.asText(),
|
||||
account,
|
||||
snapshotSlot
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
log.info(
|
||||
"Full snapshot RPC parsed. accounts={} snapshotSlot={}",
|
||||
accounts.size(),
|
||||
snapshotSlot
|
||||
);
|
||||
|
||||
return new SnapshotResult(
|
||||
snapshotSlot,
|
||||
accounts
|
||||
);
|
||||
}
|
||||
|
||||
public SignatureFetchResult getSignaturesForAddressSince(
|
||||
String address,
|
||||
String knownSignature
|
||||
) throws IOException {
|
||||
|
||||
List<SignatureRecord> signatures =
|
||||
new ArrayList<>();
|
||||
|
||||
String before =
|
||||
null;
|
||||
|
||||
boolean anchorFound =
|
||||
knownSignature == null;
|
||||
|
||||
while (true) {
|
||||
|
||||
Map<String, Object> options =
|
||||
new LinkedHashMap<>();
|
||||
|
||||
options.put(
|
||||
"limit",
|
||||
SIGNATURE_PAGE_SIZE
|
||||
);
|
||||
|
||||
options.put(
|
||||
"commitment",
|
||||
commitment
|
||||
);
|
||||
|
||||
if (before != null) {
|
||||
options.put(
|
||||
"before",
|
||||
before
|
||||
);
|
||||
}
|
||||
|
||||
Map<String, Object> payload =
|
||||
Map.of(
|
||||
"jsonrpc", "2.0",
|
||||
"id", 103,
|
||||
"method", "getSignaturesForAddress",
|
||||
"params", List.of(
|
||||
address,
|
||||
options
|
||||
)
|
||||
);
|
||||
|
||||
JsonNode root =
|
||||
executeRpc(payload);
|
||||
|
||||
JsonNode values =
|
||||
root.path("result");
|
||||
|
||||
if (!values.isArray()) {
|
||||
throw new IOException(
|
||||
"Unexpected getSignaturesForAddress response"
|
||||
);
|
||||
}
|
||||
|
||||
if (values.isEmpty()) {
|
||||
break;
|
||||
}
|
||||
|
||||
for (JsonNode item : values) {
|
||||
|
||||
String signature =
|
||||
item.path("signature")
|
||||
.asText("");
|
||||
|
||||
if (signature.isBlank()) {
|
||||
continue;
|
||||
}
|
||||
|
||||
if (knownSignature != null
|
||||
&& knownSignature.equals(signature)) {
|
||||
anchorFound = true;
|
||||
return new SignatureFetchResult(
|
||||
signatures,
|
||||
true
|
||||
);
|
||||
}
|
||||
|
||||
Long blockTime =
|
||||
item.hasNonNull("blockTime")
|
||||
? item.get("blockTime").asLong()
|
||||
: null;
|
||||
|
||||
signatures.add(
|
||||
new SignatureRecord(
|
||||
signature,
|
||||
item.path("slot")
|
||||
.asLong(-1),
|
||||
blockTime,
|
||||
item.hasNonNull("err"),
|
||||
item.path("err").toString()
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
JsonNode last =
|
||||
values.get(values.size() - 1);
|
||||
|
||||
before =
|
||||
last.path("signature")
|
||||
.asText("");
|
||||
|
||||
if (before.isBlank()
|
||||
|| values.size() < SIGNATURE_PAGE_SIZE) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
return new SignatureFetchResult(
|
||||
signatures,
|
||||
anchorFound
|
||||
);
|
||||
}
|
||||
|
||||
public JsonNode getTransactionJsonParsed(
|
||||
String signature
|
||||
) throws IOException {
|
||||
|
||||
Map<String, Object> payload =
|
||||
Map.of(
|
||||
"jsonrpc", "2.0",
|
||||
"id", 104,
|
||||
"method", "getTransaction",
|
||||
"params", List.of(
|
||||
signature,
|
||||
Map.of(
|
||||
"encoding",
|
||||
"jsonParsed",
|
||||
"commitment",
|
||||
commitment,
|
||||
"maxSupportedTransactionVersion",
|
||||
0
|
||||
)
|
||||
)
|
||||
);
|
||||
|
||||
JsonNode root =
|
||||
executeRpc(payload);
|
||||
|
||||
JsonNode result =
|
||||
root.get("result");
|
||||
|
||||
if (result == null
|
||||
|| result.isNull()) {
|
||||
return null;
|
||||
}
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
public long getCurrentSlot()
|
||||
throws IOException {
|
||||
|
||||
Map<String, Object> payload =
|
||||
Map.of(
|
||||
"jsonrpc", "2.0",
|
||||
"id", 101,
|
||||
"method", "getSlot",
|
||||
"params", List.of(
|
||||
Map.of(
|
||||
"commitment",
|
||||
commitment
|
||||
)
|
||||
)
|
||||
);
|
||||
|
||||
JsonNode root =
|
||||
executeRpc(payload);
|
||||
|
||||
long slot =
|
||||
root.path("result")
|
||||
.asLong(-1);
|
||||
|
||||
if (slot < 0) {
|
||||
throw new IOException(
|
||||
"Invalid getSlot response"
|
||||
);
|
||||
}
|
||||
|
||||
return slot;
|
||||
}
|
||||
|
||||
public AccountBatchResult getCurrentAccounts(
|
||||
Collection<String> addresses,
|
||||
long recoverySlot
|
||||
) throws IOException {
|
||||
|
||||
List<ProgramAccountUpdate> updates =
|
||||
new ArrayList<>();
|
||||
|
||||
List<String> missingAddresses =
|
||||
new ArrayList<>();
|
||||
|
||||
List<String> addressList =
|
||||
new ArrayList<>(addresses);
|
||||
|
||||
for (int offset = 0;
|
||||
offset < addressList.size();
|
||||
offset += ACCOUNT_BATCH_SIZE) {
|
||||
|
||||
int end =
|
||||
Math.min(
|
||||
offset + ACCOUNT_BATCH_SIZE,
|
||||
addressList.size()
|
||||
);
|
||||
|
||||
List<String> batch =
|
||||
addressList.subList(
|
||||
offset,
|
||||
end
|
||||
);
|
||||
|
||||
Map<String, Object> payload =
|
||||
Map.of(
|
||||
"jsonrpc", "2.0",
|
||||
"id", 105,
|
||||
"method", "getMultipleAccounts",
|
||||
"params", List.of(
|
||||
batch,
|
||||
Map.of(
|
||||
"encoding",
|
||||
"base64",
|
||||
"commitment",
|
||||
commitment
|
||||
)
|
||||
)
|
||||
);
|
||||
|
||||
JsonNode root =
|
||||
executeRpc(payload);
|
||||
|
||||
JsonNode values =
|
||||
root.path("result")
|
||||
.path("value");
|
||||
|
||||
if (!values.isArray()) {
|
||||
throw new IOException(
|
||||
"Unexpected getMultipleAccounts response"
|
||||
);
|
||||
}
|
||||
|
||||
if (values.size() != batch.size()) {
|
||||
throw new IOException(
|
||||
"Unexpected getMultipleAccounts account count"
|
||||
);
|
||||
}
|
||||
|
||||
for (int i = 0; i < batch.size(); i++) {
|
||||
|
||||
String address =
|
||||
batch.get(i);
|
||||
|
||||
JsonNode account =
|
||||
values.get(i);
|
||||
|
||||
if (account == null
|
||||
|| account.isNull()) {
|
||||
missingAddresses.add(address);
|
||||
continue;
|
||||
}
|
||||
|
||||
String owner =
|
||||
account.path("owner")
|
||||
.asText("");
|
||||
|
||||
if (!programId.equals(owner)) {
|
||||
continue;
|
||||
}
|
||||
|
||||
JsonNode data =
|
||||
account.path("data");
|
||||
|
||||
if (!data.isArray()
|
||||
|| data.isEmpty()) {
|
||||
throw new IOException(
|
||||
"Unexpected account.data format for "
|
||||
+ address
|
||||
);
|
||||
}
|
||||
|
||||
updates.add(
|
||||
parseAccount(
|
||||
address,
|
||||
account,
|
||||
recoverySlot
|
||||
)
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
return new AccountBatchResult(
|
||||
updates,
|
||||
missingAddresses
|
||||
);
|
||||
}
|
||||
|
||||
public long getFirstAvailableBlock()
|
||||
throws IOException {
|
||||
|
||||
Map<String, Object> payload =
|
||||
Map.of(
|
||||
"jsonrpc", "2.0",
|
||||
"id", 102,
|
||||
"method", "getFirstAvailableBlock"
|
||||
);
|
||||
|
||||
JsonNode root =
|
||||
executeRpc(payload);
|
||||
|
||||
return root.path("result")
|
||||
.asLong(-1);
|
||||
}
|
||||
|
||||
public List<SignatureInfo> getSignaturesAfterSlot(
|
||||
long fromSlot,
|
||||
long targetSlot
|
||||
) throws IOException {
|
||||
|
||||
List<SignatureInfo> signatures =
|
||||
new ArrayList<>();
|
||||
|
||||
String before =
|
||||
null;
|
||||
|
||||
boolean lowerBoundaryReached =
|
||||
false;
|
||||
|
||||
while (!lowerBoundaryReached) {
|
||||
|
||||
Map<String, Object> options =
|
||||
new LinkedHashMap<>();
|
||||
|
||||
options.put("limit", SIGNATURE_PAGE_SIZE);
|
||||
options.put("commitment", commitment);
|
||||
|
||||
if (before != null) {
|
||||
options.put("before", before);
|
||||
}
|
||||
|
||||
Map<String, Object> payload =
|
||||
Map.of(
|
||||
"jsonrpc", "2.0",
|
||||
"id", 106,
|
||||
"method", "getSignaturesForAddress",
|
||||
"params", List.of(
|
||||
programId,
|
||||
options
|
||||
)
|
||||
);
|
||||
|
||||
JsonNode root =
|
||||
executeRpc(payload);
|
||||
|
||||
JsonNode values =
|
||||
root.path("result");
|
||||
|
||||
if (!values.isArray()) {
|
||||
throw new IOException(
|
||||
"Unexpected getSignaturesForAddress response"
|
||||
);
|
||||
}
|
||||
|
||||
if (values.isEmpty()) {
|
||||
break;
|
||||
}
|
||||
|
||||
for (JsonNode item : values) {
|
||||
|
||||
long slot =
|
||||
item.path("slot")
|
||||
.asLong(-1);
|
||||
|
||||
if (slot < 0) {
|
||||
continue;
|
||||
}
|
||||
|
||||
if (slot <= fromSlot) {
|
||||
lowerBoundaryReached = true;
|
||||
break;
|
||||
}
|
||||
|
||||
if (slot > targetSlot) {
|
||||
continue;
|
||||
}
|
||||
|
||||
String signature =
|
||||
item.path("signature")
|
||||
.asText("");
|
||||
|
||||
if (signature.isBlank()) {
|
||||
continue;
|
||||
}
|
||||
|
||||
JsonNode error =
|
||||
item.get("err");
|
||||
|
||||
if (error != null
|
||||
&& !error.isNull()) {
|
||||
continue;
|
||||
}
|
||||
|
||||
signatures.add(
|
||||
new SignatureInfo(
|
||||
signature,
|
||||
slot
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
JsonNode last =
|
||||
values.get(values.size() - 1);
|
||||
|
||||
before =
|
||||
last.path("signature")
|
||||
.asText("");
|
||||
|
||||
if (before.isBlank()
|
||||
|| values.size() < SIGNATURE_PAGE_SIZE) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
return signatures;
|
||||
}
|
||||
|
||||
public Set<String> getTouchedAddresses(
|
||||
List<SignatureInfo> signatures
|
||||
) throws IOException {
|
||||
|
||||
Set<String> addresses =
|
||||
new LinkedHashSet<>();
|
||||
|
||||
for (SignatureInfo signatureInfo : signatures) {
|
||||
|
||||
JsonNode transaction =
|
||||
getTransactionJsonParsed(
|
||||
signatureInfo.signature()
|
||||
);
|
||||
|
||||
if (transaction == null
|
||||
|| transaction.isNull()) {
|
||||
throw new IOException(
|
||||
"Transaction unavailable during recovery: "
|
||||
+ signatureInfo.signature()
|
||||
);
|
||||
}
|
||||
|
||||
JsonNode accountKeys =
|
||||
transaction.path("transaction")
|
||||
.path("message")
|
||||
.path("accountKeys");
|
||||
|
||||
if (!accountKeys.isArray()) {
|
||||
throw new IOException(
|
||||
"Missing accountKeys for transaction: "
|
||||
+ signatureInfo.signature()
|
||||
);
|
||||
}
|
||||
|
||||
for (JsonNode keyNode : accountKeys) {
|
||||
|
||||
String pubkey;
|
||||
|
||||
if (keyNode.isTextual()) {
|
||||
pubkey = keyNode.asText();
|
||||
} else {
|
||||
pubkey = keyNode.path("pubkey")
|
||||
.asText("");
|
||||
}
|
||||
|
||||
if (!pubkey.isBlank()) {
|
||||
addresses.add(pubkey);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
addresses.remove(programId);
|
||||
return addresses;
|
||||
}
|
||||
|
||||
private ProgramAccountUpdate parseAccount(
|
||||
String address,
|
||||
JsonNode account,
|
||||
long slot
|
||||
) {
|
||||
|
||||
JsonNode data =
|
||||
account.path("data");
|
||||
|
||||
return new ProgramAccountUpdate(
|
||||
address,
|
||||
account.path("owner")
|
||||
.asText(),
|
||||
account.path("lamports")
|
||||
.asLong(),
|
||||
slot,
|
||||
data.get(0)
|
||||
.asText(),
|
||||
account.path("executable")
|
||||
.asBoolean(false),
|
||||
account.hasNonNull("rentEpoch")
|
||||
? account.get("rentEpoch").asLong()
|
||||
: null
|
||||
);
|
||||
}
|
||||
|
||||
private JsonNode executeRpc(
|
||||
Map<String, Object> payload
|
||||
) throws IOException {
|
||||
|
||||
Request request =
|
||||
new Request.Builder()
|
||||
.url(rpcUrl)
|
||||
.post(
|
||||
RequestBody.create(
|
||||
mapper.writeValueAsBytes(payload),
|
||||
JSON
|
||||
)
|
||||
)
|
||||
.build();
|
||||
|
||||
try (Response response =
|
||||
httpClient.newCall(request).execute()) {
|
||||
|
||||
if (!response.isSuccessful()) {
|
||||
throw new IOException(
|
||||
"Solana RPC HTTP error: " + response.code()
|
||||
);
|
||||
}
|
||||
|
||||
ResponseBody body =
|
||||
response.body();
|
||||
|
||||
if (body == null) {
|
||||
throw new IOException(
|
||||
"Solana RPC returned empty body"
|
||||
);
|
||||
}
|
||||
|
||||
JsonNode root =
|
||||
mapper.readTree(
|
||||
body.string()
|
||||
);
|
||||
|
||||
if (root.has("error")) {
|
||||
throw new IOException(
|
||||
"Solana RPC error: " + root.get("error")
|
||||
);
|
||||
}
|
||||
|
||||
return root;
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
|
||||
httpClient.dispatcher()
|
||||
.executorService()
|
||||
.shutdown();
|
||||
httpClient.connectionPool()
|
||||
.evictAll();
|
||||
}
|
||||
|
||||
public record SignatureRecord(
|
||||
String signature,
|
||||
long slot,
|
||||
Long blockTime,
|
||||
boolean failed,
|
||||
String errorJson
|
||||
) {
|
||||
}
|
||||
|
||||
public record SignatureFetchResult(
|
||||
List<SignatureRecord> signatures,
|
||||
boolean anchorFound
|
||||
) {
|
||||
}
|
||||
|
||||
public record AccountBatchResult(
|
||||
List<ProgramAccountUpdate> updates,
|
||||
List<String> missingAddresses
|
||||
) {
|
||||
}
|
||||
|
||||
public record SignatureInfo(
|
||||
String signature,
|
||||
long slot
|
||||
) {
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,542 @@
|
||||
package sync.source.rpc;
|
||||
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import okhttp3.*;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import sync.model.ProgramAccountUpdate;
|
||||
import sync.source.AccountUpdateListener;
|
||||
import sync.source.ConnectionListener;
|
||||
|
||||
import java.time.Duration;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.*;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
public final class SolanaWebSocketClient
|
||||
extends WebSocketListener
|
||||
implements AutoCloseable {
|
||||
|
||||
private static final Logger log =
|
||||
LoggerFactory.getLogger(
|
||||
SolanaWebSocketClient.class
|
||||
);
|
||||
|
||||
private final String websocketUrl;
|
||||
private final String programId;
|
||||
private final String commitment;
|
||||
|
||||
private final ObjectMapper mapper =
|
||||
new ObjectMapper();
|
||||
|
||||
private final OkHttpClient httpClient;
|
||||
|
||||
private final ScheduledExecutorService scheduler;
|
||||
|
||||
private final AtomicBoolean closed =
|
||||
new AtomicBoolean(false);
|
||||
|
||||
private final AtomicInteger reconnectAttempt =
|
||||
new AtomicInteger(0);
|
||||
|
||||
private final AtomicBoolean everSubscribed =
|
||||
new AtomicBoolean(false);
|
||||
|
||||
private volatile WebSocket webSocket;
|
||||
|
||||
private volatile AccountUpdateListener
|
||||
accountUpdateListener;
|
||||
|
||||
private volatile ConnectionListener
|
||||
connectionListener;
|
||||
|
||||
public SolanaWebSocketClient(
|
||||
String websocketUrl,
|
||||
String programId,
|
||||
String commitment
|
||||
) {
|
||||
this.websocketUrl = websocketUrl;
|
||||
this.programId = programId;
|
||||
this.commitment = commitment;
|
||||
|
||||
this.httpClient =
|
||||
new OkHttpClient.Builder()
|
||||
.readTimeout(
|
||||
Duration.ZERO
|
||||
)
|
||||
.pingInterval(
|
||||
Duration.ofSeconds(20)
|
||||
)
|
||||
.build();
|
||||
|
||||
this.scheduler =
|
||||
Executors
|
||||
.newSingleThreadScheduledExecutor(
|
||||
runnable -> {
|
||||
|
||||
Thread thread =
|
||||
new Thread(
|
||||
runnable,
|
||||
"solana-rpc-reconnect"
|
||||
);
|
||||
|
||||
thread.setDaemon(
|
||||
true
|
||||
);
|
||||
|
||||
return thread;
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
public void start(
|
||||
AccountUpdateListener accountUpdateListener,
|
||||
ConnectionListener connectionListener
|
||||
) {
|
||||
|
||||
this.accountUpdateListener =
|
||||
accountUpdateListener;
|
||||
|
||||
this.connectionListener =
|
||||
connectionListener;
|
||||
|
||||
connect();
|
||||
}
|
||||
|
||||
private void connect() {
|
||||
|
||||
if (closed.get()) {
|
||||
return;
|
||||
}
|
||||
|
||||
log.info(
|
||||
"Connecting to Solana WebSocket: {}",
|
||||
websocketUrl
|
||||
);
|
||||
|
||||
Request request =
|
||||
new Request.Builder()
|
||||
.url(
|
||||
websocketUrl
|
||||
)
|
||||
.build();
|
||||
|
||||
this.webSocket =
|
||||
httpClient.newWebSocket(
|
||||
request,
|
||||
this
|
||||
);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onOpen(
|
||||
WebSocket webSocket,
|
||||
Response response
|
||||
) {
|
||||
|
||||
reconnectAttempt.set(
|
||||
0
|
||||
);
|
||||
|
||||
log.info(
|
||||
"WebSocket connected"
|
||||
);
|
||||
|
||||
sendProgramSubscribe(
|
||||
webSocket
|
||||
);
|
||||
}
|
||||
|
||||
private void sendProgramSubscribe(
|
||||
WebSocket webSocket
|
||||
) {
|
||||
|
||||
try {
|
||||
|
||||
Map<String, Object> request =
|
||||
Map.of(
|
||||
"jsonrpc",
|
||||
"2.0",
|
||||
"id",
|
||||
1,
|
||||
"method",
|
||||
"programSubscribe",
|
||||
"params",
|
||||
List.of(
|
||||
programId,
|
||||
Map.of(
|
||||
"encoding",
|
||||
"base64",
|
||||
"commitment",
|
||||
commitment
|
||||
)
|
||||
)
|
||||
);
|
||||
|
||||
String payload =
|
||||
mapper.writeValueAsString(
|
||||
request
|
||||
);
|
||||
|
||||
if (!webSocket.send(
|
||||
payload
|
||||
)) {
|
||||
throw new IllegalStateException(
|
||||
"WebSocket rejected subscription request"
|
||||
);
|
||||
}
|
||||
|
||||
log.info(
|
||||
"Subscription request sent for program: {}",
|
||||
programId
|
||||
);
|
||||
|
||||
} catch (Exception exception) {
|
||||
|
||||
log.error(
|
||||
"Failed to send subscription request",
|
||||
exception
|
||||
);
|
||||
|
||||
webSocket.cancel();
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onMessage(
|
||||
WebSocket webSocket,
|
||||
String text
|
||||
) {
|
||||
|
||||
try {
|
||||
|
||||
JsonNode root =
|
||||
mapper.readTree(
|
||||
text
|
||||
);
|
||||
|
||||
if (
|
||||
root.has("id")
|
||||
&& root.has("result")
|
||||
&& root.get("id").asInt() == 1
|
||||
) {
|
||||
|
||||
log.info(
|
||||
"Subscribed successfully. subscriptionId={}",
|
||||
root.get("result").asText()
|
||||
);
|
||||
|
||||
boolean firstConnection =
|
||||
everSubscribed
|
||||
.compareAndSet(
|
||||
false,
|
||||
true
|
||||
);
|
||||
|
||||
ConnectionListener listener =
|
||||
connectionListener;
|
||||
|
||||
if (listener != null) {
|
||||
listener.onConnected(
|
||||
firstConnection
|
||||
);
|
||||
}
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
if (root.has("error")) {
|
||||
|
||||
log.error(
|
||||
"Solana websocket RPC error: {}",
|
||||
root.get("error")
|
||||
);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
if (!"programNotification"
|
||||
.equals(
|
||||
root
|
||||
.path("method")
|
||||
.asText()
|
||||
)) {
|
||||
return;
|
||||
}
|
||||
|
||||
ProgramAccountUpdate update =
|
||||
parseProgramNotification(
|
||||
root
|
||||
);
|
||||
|
||||
log.debug(
|
||||
"Account update. address={} slot={} base64Chars={}",
|
||||
update.address(),
|
||||
update.slot(),
|
||||
update.dataBase64().length()
|
||||
);
|
||||
|
||||
AccountUpdateListener listener =
|
||||
accountUpdateListener;
|
||||
|
||||
if (listener != null) {
|
||||
listener.onAccountUpdate(
|
||||
update
|
||||
);
|
||||
}
|
||||
|
||||
} catch (Exception exception) {
|
||||
|
||||
log.error(
|
||||
"Failed to process WebSocket message. raw={}",
|
||||
text,
|
||||
exception
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
private ProgramAccountUpdate parseProgramNotification(
|
||||
JsonNode root
|
||||
) {
|
||||
|
||||
JsonNode result =
|
||||
root
|
||||
.path("params")
|
||||
.path("result");
|
||||
|
||||
JsonNode context =
|
||||
result.path("context");
|
||||
|
||||
JsonNode value =
|
||||
result.path("value");
|
||||
|
||||
JsonNode account =
|
||||
value.path("account");
|
||||
|
||||
JsonNode data =
|
||||
account.path("data");
|
||||
|
||||
if (!data.isArray()
|
||||
|| data.isEmpty()) {
|
||||
throw new IllegalArgumentException(
|
||||
"Unexpected account.data format"
|
||||
);
|
||||
}
|
||||
|
||||
return new ProgramAccountUpdate(
|
||||
requiredText(
|
||||
value,
|
||||
"pubkey"
|
||||
),
|
||||
requiredText(
|
||||
account,
|
||||
"owner"
|
||||
),
|
||||
requiredLong(
|
||||
account,
|
||||
"lamports"
|
||||
),
|
||||
requiredLong(
|
||||
context,
|
||||
"slot"
|
||||
),
|
||||
data.get(0).asText(),
|
||||
account.path(
|
||||
"executable"
|
||||
).asBoolean(false),
|
||||
account.hasNonNull(
|
||||
"rentEpoch"
|
||||
)
|
||||
? account
|
||||
.get("rentEpoch")
|
||||
.asLong()
|
||||
: null
|
||||
);
|
||||
}
|
||||
|
||||
private String requiredText(
|
||||
JsonNode node,
|
||||
String field
|
||||
) {
|
||||
|
||||
JsonNode value =
|
||||
node.get(field);
|
||||
|
||||
if (value == null
|
||||
|| value.isNull()
|
||||
|| value.asText().isBlank()) {
|
||||
throw new IllegalArgumentException(
|
||||
"Missing field: "
|
||||
+ field
|
||||
);
|
||||
}
|
||||
|
||||
return value.asText();
|
||||
}
|
||||
|
||||
private long requiredLong(
|
||||
JsonNode node,
|
||||
String field
|
||||
) {
|
||||
|
||||
JsonNode value =
|
||||
node.get(field);
|
||||
|
||||
if (value == null
|
||||
|| !value.isNumber()) {
|
||||
throw new IllegalArgumentException(
|
||||
"Missing numeric field: "
|
||||
+ field
|
||||
);
|
||||
}
|
||||
|
||||
return value.asLong();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onClosed(
|
||||
WebSocket webSocket,
|
||||
int code,
|
||||
String reason
|
||||
) {
|
||||
|
||||
log.warn(
|
||||
"WebSocket closed. code={} reason={}",
|
||||
code,
|
||||
reason
|
||||
);
|
||||
|
||||
notifyDisconnected(
|
||||
new IllegalStateException(
|
||||
"WebSocket closed: "
|
||||
+ reason
|
||||
)
|
||||
);
|
||||
|
||||
scheduleReconnect();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onFailure(
|
||||
WebSocket webSocket,
|
||||
Throwable throwable,
|
||||
Response response
|
||||
) {
|
||||
|
||||
log.error(
|
||||
"WebSocket failure",
|
||||
throwable
|
||||
);
|
||||
|
||||
if (response != null) {
|
||||
log.error(
|
||||
"WebSocket HTTP status: {}",
|
||||
response.code()
|
||||
);
|
||||
}
|
||||
|
||||
notifyDisconnected(
|
||||
throwable
|
||||
);
|
||||
|
||||
scheduleReconnect();
|
||||
}
|
||||
|
||||
private void notifyDisconnected(
|
||||
Throwable cause
|
||||
) {
|
||||
|
||||
ConnectionListener listener =
|
||||
connectionListener;
|
||||
|
||||
if (listener != null) {
|
||||
|
||||
try {
|
||||
listener.onDisconnected(
|
||||
cause
|
||||
);
|
||||
} catch (Exception exception) {
|
||||
log.error(
|
||||
"Connection listener failed",
|
||||
exception
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void scheduleReconnect() {
|
||||
|
||||
if (closed.get()) {
|
||||
return;
|
||||
}
|
||||
|
||||
int attempt =
|
||||
reconnectAttempt
|
||||
.incrementAndGet();
|
||||
|
||||
long delaySeconds =
|
||||
switch (
|
||||
Math.min(
|
||||
attempt,
|
||||
5
|
||||
)
|
||||
) {
|
||||
case 1 -> 1;
|
||||
case 2 -> 2;
|
||||
case 3 -> 5;
|
||||
case 4 -> 10;
|
||||
default -> 30;
|
||||
};
|
||||
|
||||
log.warn(
|
||||
"Reconnect scheduled in {} seconds (attempt {})",
|
||||
delaySeconds,
|
||||
attempt
|
||||
);
|
||||
|
||||
scheduler.schedule(
|
||||
this::connect,
|
||||
delaySeconds,
|
||||
TimeUnit.SECONDS
|
||||
);
|
||||
}
|
||||
|
||||
public void stop() {
|
||||
close();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
|
||||
if (!closed.compareAndSet(
|
||||
false,
|
||||
true
|
||||
)) {
|
||||
return;
|
||||
}
|
||||
|
||||
WebSocket socket =
|
||||
this.webSocket;
|
||||
|
||||
if (socket != null) {
|
||||
socket.close(
|
||||
1000,
|
||||
"Application shutdown"
|
||||
);
|
||||
}
|
||||
|
||||
scheduler.shutdownNow();
|
||||
|
||||
httpClient
|
||||
.dispatcher()
|
||||
.executorService()
|
||||
.shutdown();
|
||||
|
||||
httpClient
|
||||
.connectionPool()
|
||||
.evictAll();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,224 @@
|
||||
package sync.util;
|
||||
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.util.Arrays;
|
||||
|
||||
public final class Base58Util {
|
||||
|
||||
private static final char[] ALPHABET =
|
||||
"123456789ABCDEFGHJKLMNPQRSTUVWXYZabcdefghijkmnopqrstuvwxyz"
|
||||
.toCharArray();
|
||||
|
||||
private static final int[] INDEXES =
|
||||
new int[128];
|
||||
|
||||
static {
|
||||
Arrays.fill(
|
||||
INDEXES,
|
||||
-1
|
||||
);
|
||||
|
||||
for (int i = 0; i < ALPHABET.length; i++) {
|
||||
INDEXES[ALPHABET[i]] = i;
|
||||
}
|
||||
}
|
||||
|
||||
private Base58Util() {
|
||||
}
|
||||
|
||||
public static String encode(
|
||||
byte[] input
|
||||
) {
|
||||
|
||||
if (input == null
|
||||
|| input.length == 0) {
|
||||
return "";
|
||||
}
|
||||
|
||||
byte[] copy =
|
||||
Arrays.copyOf(
|
||||
input,
|
||||
input.length
|
||||
);
|
||||
|
||||
int zeros =
|
||||
0;
|
||||
|
||||
while (zeros < copy.length
|
||||
&& copy[zeros] == 0) {
|
||||
zeros++;
|
||||
}
|
||||
|
||||
byte[] encoded =
|
||||
new byte[copy.length * 2];
|
||||
|
||||
int outputStart =
|
||||
encoded.length;
|
||||
|
||||
int inputStart =
|
||||
zeros;
|
||||
|
||||
while (inputStart < copy.length) {
|
||||
|
||||
int remainder =
|
||||
divMod58(
|
||||
copy,
|
||||
inputStart
|
||||
);
|
||||
|
||||
if (copy[inputStart] == 0) {
|
||||
inputStart++;
|
||||
}
|
||||
|
||||
encoded[--outputStart] =
|
||||
(byte) ALPHABET[remainder];
|
||||
}
|
||||
|
||||
while (outputStart < encoded.length
|
||||
&& encoded[outputStart] == ALPHABET[0]) {
|
||||
outputStart++;
|
||||
}
|
||||
|
||||
while (--zeros >= 0) {
|
||||
encoded[--outputStart] =
|
||||
(byte) ALPHABET[0];
|
||||
}
|
||||
|
||||
return new String(
|
||||
encoded,
|
||||
outputStart,
|
||||
encoded.length - outputStart,
|
||||
StandardCharsets.US_ASCII
|
||||
);
|
||||
}
|
||||
|
||||
public static byte[] decode(
|
||||
String input
|
||||
) {
|
||||
|
||||
if (input == null
|
||||
|| input.isBlank()) {
|
||||
return new byte[0];
|
||||
}
|
||||
|
||||
char[] chars =
|
||||
input.trim()
|
||||
.toCharArray();
|
||||
|
||||
byte[] input58 =
|
||||
new byte[chars.length];
|
||||
|
||||
for (int i = 0; i < chars.length; i++) {
|
||||
|
||||
char c = chars[i];
|
||||
|
||||
if (c >= INDEXES.length
|
||||
|| INDEXES[c] < 0) {
|
||||
throw new IllegalArgumentException(
|
||||
"Invalid Base58 character: " + c
|
||||
);
|
||||
}
|
||||
|
||||
input58[i] =
|
||||
(byte) INDEXES[c];
|
||||
}
|
||||
|
||||
int zeros =
|
||||
0;
|
||||
|
||||
while (zeros < input58.length
|
||||
&& input58[zeros] == 0) {
|
||||
zeros++;
|
||||
}
|
||||
|
||||
byte[] decoded =
|
||||
new byte[chars.length];
|
||||
|
||||
int outputStart =
|
||||
decoded.length;
|
||||
|
||||
int inputStart =
|
||||
zeros;
|
||||
|
||||
while (inputStart < input58.length) {
|
||||
|
||||
int remainder =
|
||||
divMod256(
|
||||
input58,
|
||||
inputStart
|
||||
);
|
||||
|
||||
if (input58[inputStart] == 0) {
|
||||
inputStart++;
|
||||
}
|
||||
|
||||
decoded[--outputStart] =
|
||||
(byte) remainder;
|
||||
}
|
||||
|
||||
while (outputStart < decoded.length
|
||||
&& decoded[outputStart] == 0) {
|
||||
outputStart++;
|
||||
}
|
||||
|
||||
return Arrays.copyOfRange(
|
||||
decoded,
|
||||
outputStart - zeros,
|
||||
decoded.length
|
||||
);
|
||||
}
|
||||
|
||||
private static int divMod58(
|
||||
byte[] number,
|
||||
int startAt
|
||||
) {
|
||||
|
||||
int remainder =
|
||||
0;
|
||||
|
||||
for (int i = startAt; i < number.length; i++) {
|
||||
|
||||
int digit256 =
|
||||
number[i] & 0xFF;
|
||||
|
||||
int temp =
|
||||
remainder * 256
|
||||
+ digit256;
|
||||
|
||||
number[i] =
|
||||
(byte) (temp / 58);
|
||||
|
||||
remainder =
|
||||
temp % 58;
|
||||
}
|
||||
|
||||
return remainder;
|
||||
}
|
||||
|
||||
private static int divMod256(
|
||||
byte[] number58,
|
||||
int startAt
|
||||
) {
|
||||
|
||||
int remainder =
|
||||
0;
|
||||
|
||||
for (int i = startAt; i < number58.length; i++) {
|
||||
|
||||
int digit58 =
|
||||
number58[i] & 0xFF;
|
||||
|
||||
int temp =
|
||||
remainder * 58
|
||||
+ digit58;
|
||||
|
||||
number58[i] =
|
||||
(byte) (temp / 256);
|
||||
|
||||
remainder =
|
||||
temp % 256;
|
||||
}
|
||||
|
||||
return remainder;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,95 @@
|
||||
package sync.util;
|
||||
|
||||
import org.bouncycastle.jcajce.provider.digest.SHA256;
|
||||
import org.bouncycastle.math.ec.rfc8032.Ed25519;
|
||||
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
|
||||
public final class SolanaPdaUtil {
|
||||
|
||||
private static final byte[] PROGRAM_DERIVED_ADDRESS_DOMAIN =
|
||||
"ProgramDerivedAddress"
|
||||
.getBytes(StandardCharsets.UTF_8);
|
||||
|
||||
private SolanaPdaUtil() {
|
||||
}
|
||||
|
||||
public static String findProgramAddress(
|
||||
List<byte[]> seeds,
|
||||
String programIdBase58
|
||||
) {
|
||||
|
||||
byte[] programId =
|
||||
Base58Util.decode(
|
||||
programIdBase58
|
||||
);
|
||||
|
||||
if (programId.length != 32) {
|
||||
throw new IllegalArgumentException(
|
||||
"Program id must decode to 32 bytes"
|
||||
);
|
||||
}
|
||||
|
||||
for (int bump = 255; bump >= 0; bump--) {
|
||||
|
||||
List<byte[]> attemptSeeds =
|
||||
new ArrayList<>(
|
||||
seeds
|
||||
);
|
||||
|
||||
attemptSeeds.add(
|
||||
new byte[]{(byte) bump}
|
||||
);
|
||||
|
||||
byte[] address =
|
||||
createProgramAddress(
|
||||
attemptSeeds,
|
||||
programId
|
||||
);
|
||||
|
||||
if (!isOnCurve(address)) {
|
||||
return Base58Util.encode(
|
||||
address
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
throw new IllegalStateException(
|
||||
"Unable to find viable program address"
|
||||
);
|
||||
}
|
||||
|
||||
private static byte[] createProgramAddress(
|
||||
List<byte[]> seeds,
|
||||
byte[] programId
|
||||
) {
|
||||
|
||||
SHA256.Digest digest =
|
||||
new SHA256.Digest();
|
||||
|
||||
for (byte[] seed : seeds) {
|
||||
digest.update(seed);
|
||||
}
|
||||
|
||||
digest.update(programId);
|
||||
digest.update(PROGRAM_DERIVED_ADDRESS_DOMAIN);
|
||||
|
||||
return digest.digest();
|
||||
}
|
||||
|
||||
private static boolean isOnCurve(
|
||||
byte[] publicKey
|
||||
) {
|
||||
|
||||
if (publicKey.length != 32) {
|
||||
return false;
|
||||
}
|
||||
|
||||
return Ed25519.validatePublicKeyFull(
|
||||
publicKey,
|
||||
0
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,85 @@
|
||||
package server.sync;
|
||||
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import sync.service.SolanaUsersSyncService;
|
||||
import utils.config.AppConfig;
|
||||
|
||||
public final class SolanaUsersSyncStartupService {
|
||||
|
||||
private static final Logger log =
|
||||
LoggerFactory.getLogger(
|
||||
SolanaUsersSyncStartupService.class
|
||||
);
|
||||
|
||||
private static volatile SolanaUsersSyncService service;
|
||||
|
||||
private SolanaUsersSyncStartupService() {
|
||||
}
|
||||
|
||||
public static synchronized void startOrThrow()
|
||||
throws Exception {
|
||||
|
||||
if (service != null) {
|
||||
return;
|
||||
}
|
||||
|
||||
AppConfig serverConfig =
|
||||
AppConfig.getInstance();
|
||||
|
||||
if (!sync.config.AppConfig.isEnabled(serverConfig)) {
|
||||
log.info(
|
||||
"Solana users sync is disabled. Param {} is false or absent.",
|
||||
sync.config.AppConfig.ENABLED_KEY
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
sync.config.AppConfig syncConfig =
|
||||
sync.config.AppConfig.fromServerConfig(
|
||||
serverConfig
|
||||
);
|
||||
|
||||
service =
|
||||
new SolanaUsersSyncService(
|
||||
syncConfig
|
||||
);
|
||||
|
||||
Runtime.getRuntime()
|
||||
.addShutdownHook(
|
||||
new Thread(
|
||||
SolanaUsersSyncStartupService::closeQuietly,
|
||||
"solana-users-sync-shutdown"
|
||||
)
|
||||
);
|
||||
|
||||
log.info(
|
||||
"Starting Solana users sync before remaining server startup..."
|
||||
);
|
||||
|
||||
service.start();
|
||||
service.awaitReady();
|
||||
|
||||
log.info(
|
||||
"Solana users sync is READY. Server startup may continue."
|
||||
);
|
||||
}
|
||||
|
||||
public static synchronized void closeQuietly() {
|
||||
|
||||
if (service == null) {
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
service.close();
|
||||
} catch (Exception exception) {
|
||||
log.warn(
|
||||
"Failed to close Solana users sync service cleanly",
|
||||
exception
|
||||
);
|
||||
} finally {
|
||||
service = null;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -8,6 +8,7 @@ import org.slf4j.LoggerFactory;
|
||||
import server.debug.DebugApiConfigurator;
|
||||
import server.sync.BlockchainResyncRecoveryOnStartup;
|
||||
import server.sync.PeriodicBlockchainSyncService;
|
||||
import server.sync.SolanaUsersSyncStartupService;
|
||||
import server.sync.SyncServersBootstrapService;
|
||||
import utils.config.AppConfig;
|
||||
|
||||
@@ -49,6 +50,11 @@ public final class WsServer {
|
||||
throw e;
|
||||
}
|
||||
|
||||
// ============================================================
|
||||
// 1.0) Синхронизация пользовательских Solana PDA
|
||||
// ============================================================
|
||||
SolanaUsersSyncStartupService.startOrThrow();
|
||||
|
||||
// ============================================================
|
||||
// 1) Настройки порта
|
||||
// ============================================================
|
||||
|
||||
@@ -3,6 +3,14 @@ db.path=data/shine.sqlite
|
||||
server.SHiNE.login=shineupme
|
||||
solana.cluster=mainnet-beta
|
||||
solana.rpcUrl=https://api.mainnet-beta.solana.com
|
||||
solana.users.sync.enabled=false
|
||||
solana.users.sync.rpcUrl=
|
||||
solana.users.sync.wsUrl=
|
||||
solana.users.sync.programId=SHiNEPr1APdAgNBteUyBXcNovaHctpSjUu8oH2ZJdN6
|
||||
solana.users.sync.databaseUrl=
|
||||
solana.users.sync.dbUser=
|
||||
solana.users.sync.dbPassword=
|
||||
solana.users.sync.pollIntervalSeconds=300
|
||||
|
||||
# ------------------------------------------------------------
|
||||
# Межсерверная синхронизация: как создавать локальную запись пользователя,
|
||||
@@ -44,7 +52,7 @@ webpush.vapid.subject=mailto:admin@shine.local
|
||||
# Тогда сервер будет выдавать временный username/password (TTL).
|
||||
# ------------------------------------------------------------
|
||||
call.ice.stun.urls=stun:stun.l.google.com:19302
|
||||
call.ice.turn.urls=turn:185.229.109.118:3478?transport=udp,turn:185.229.109.118:3478?transport=tcp
|
||||
call.ice.turn.urls=turn:turn1.shineup.me:3478?transport=udp,turn:turn1.shineup.me:3478?transport=tcp
|
||||
call.ice.turn.ttlSec=600
|
||||
call.ice.turn.userPrefix=shine
|
||||
call.ice.turn.sharedSecret=
|
||||
@@ -58,12 +66,24 @@ call.ice.turn.password=
|
||||
# Каждый блок описывает один TURN-узел. Новые узлы добавляются по индексу.
|
||||
# Приоритет авторизации на узел: sharedSecret -> статические username/password.
|
||||
# ------------------------------------------------------------
|
||||
call.ice.turn.servers.1.id=shineup-main-185
|
||||
call.ice.turn.servers.1.urls=turn:185.229.109.118:3478?transport=udp,turn:185.229.109.118:3478?transport=tcp
|
||||
call.ice.turn.servers.1.sharedSecret=def6d444734d380d2f67a9d345b1debf985eaba0973c343e392c060d97c30106
|
||||
call.ice.turn.servers.1.id=turn1
|
||||
call.ice.turn.servers.1.urls=turn:turn1.shineup.me:3478?transport=udp,turn:turn1.shineup.me:3478?transport=tcp
|
||||
call.ice.turn.servers.1.sharedSecret=
|
||||
call.ice.turn.servers.1.username=
|
||||
call.ice.turn.servers.1.password=
|
||||
|
||||
call.ice.turn.servers.2.id=turn2
|
||||
call.ice.turn.servers.2.urls=turn:turn2.shineup.me:3478?transport=udp,turn:turn2.shineup.me:3478?transport=tcp
|
||||
call.ice.turn.servers.2.sharedSecret=
|
||||
call.ice.turn.servers.2.username=
|
||||
call.ice.turn.servers.2.password=
|
||||
|
||||
call.ice.turn.servers.3.id=turn3
|
||||
call.ice.turn.servers.3.urls=turn:turn3.shineup.me:3478?transport=udp,turn:turn3.shineup.me:3478?transport=tcp
|
||||
call.ice.turn.servers.3.sharedSecret=
|
||||
call.ice.turn.servers.3.username=
|
||||
call.ice.turn.servers.3.password=
|
||||
|
||||
# ------------------------------------------------------------
|
||||
# Временные debug HTTP API для тестирования соединений
|
||||
# true - endpoint'ы /debug/ws/* включены (только при наличии .debug-token)
|
||||
|
||||
@@ -1,220 +0,0 @@
|
||||
# Задача 01: Доработка вкладки «Каналы» (UI + API)
|
||||
|
||||
## Кратко и по делу
|
||||
Нужно довести вторую вкладку «Каналы» до полностью рабочего состояния на реальных данных сервера.
|
||||
|
||||
Что должно работать:
|
||||
- список каналов;
|
||||
- вход в канал и чтение сообщений;
|
||||
- вход в тред сообщения (история/ветка);
|
||||
- ответ на сообщение;
|
||||
- лайк/снятие лайка;
|
||||
- подписка на пользователя;
|
||||
- подписка на канал;
|
||||
- видимое имя канала в формате `имя_пользователя/имя_канала`.
|
||||
|
||||
Запись любых новых сущностей делается через `AddBlock` с подписью на клиенте.
|
||||
Чтение делается через 3 API:
|
||||
- `ListSubscriptionsFeed`
|
||||
- `GetChannelMessages`
|
||||
- `GetMessageThread`
|
||||
|
||||
Техническая особенность (оставляем как есть):
|
||||
- на экране каналов индикатор непрочитанного = общее число сообщений канала.
|
||||
|
||||
---
|
||||
|
||||
## Подробное ТЗ
|
||||
|
||||
### 1. Цель
|
||||
Сделать рабочий каналовый сценарий «от списка до треда», где чтение строится на RPC API, а запись действий пользователя — только через `AddBlock`.
|
||||
|
||||
### 2. Что уже есть в проекте
|
||||
|
||||
#### 2.1 UI (частично)
|
||||
- Есть страницы:
|
||||
- `channels-list`
|
||||
- `channel-view`
|
||||
- `add-channel-view`
|
||||
- Есть запросы чтения в клиенте:
|
||||
- `authService.listSubscriptionsFeed(...)`
|
||||
- `authService.getChannelMessages(...)`
|
||||
- `authService.getMessageThread(...)`
|
||||
- Есть fallback на mock-данные при ошибках сервера.
|
||||
|
||||
#### 2.2 API/сервер (уже реализованы)
|
||||
- `ListSubscriptionsFeed`
|
||||
- `GetChannelMessages`
|
||||
- `GetMessageThread`
|
||||
- `AddBlock`
|
||||
|
||||
#### 2.3 Тесты
|
||||
- Есть интеграционный тест API каналов: `IT_06_ChannelsApi`.
|
||||
- Есть тесты генерации блоков каналов/связей: `IT_03_AddBlock_NoAuth`.
|
||||
- Формат `AddBlock` и его сборка/подпись описаны в `AddBlockSender`.
|
||||
|
||||
### 3. Проблемы текущей реализации (что надо закрыть)
|
||||
- Кнопки «подписаться на человека/канал» в списке каналов сейчас UI-only (модалка без реальной записи через `AddBlock`).
|
||||
- `add-channel-view` пока не создает канал на сервере через `AddBlock` (`CreateChannelBody`), только делает `navigate`.
|
||||
- `channel-view` добавляет пост локально (в память), а не отправляет блок `TEXT_POST` через `AddBlock`.
|
||||
- Нет полноценного экрана треда сообщения с реальными `GetMessageThread` и действиями `ответить/лайк/убрать лайк` через блоки.
|
||||
- Нет гарантированного отображения канала в требуемом формате `ownerLogin/channelName`.
|
||||
|
||||
### 4. Функциональные требования
|
||||
|
||||
#### 4.1 Список каналов
|
||||
На вкладке «Каналы» отображать 3 группы:
|
||||
- Мои каналы
|
||||
- Каналы пользователей, на кого я подписан
|
||||
- Каналы, на которые я подписан
|
||||
|
||||
Источник данных: `ListSubscriptionsFeed`.
|
||||
|
||||
Каждый канал показывать в формате:
|
||||
- `ownerLogin/channelName`
|
||||
|
||||
#### 4.2 Открытие канала
|
||||
При входе в канал:
|
||||
- загрузить сообщения через `GetChannelMessages`;
|
||||
- показать список сообщений в хронологическом порядке (по текущему параметру `sort`);
|
||||
- оставить техническую особенность непрочитанных как есть.
|
||||
|
||||
#### 4.3 Открытие треда сообщения
|
||||
При клике на сообщение:
|
||||
- загрузить тред через `GetMessageThread`;
|
||||
- показать `ancestors`, `focus`, `descendants`;
|
||||
- из треда должны быть доступны действия:
|
||||
- «Ответить»
|
||||
- «Лайк»
|
||||
- «Убрать лайк»
|
||||
|
||||
Запись действий — только `AddBlock`.
|
||||
|
||||
#### 4.4 Создание канала
|
||||
В `add-channel-view` кнопка «Создать» должна:
|
||||
- отправлять `AddBlock` с телом `CreateChannelBody`;
|
||||
- после успеха возвращать к списку каналов и обновлять его.
|
||||
|
||||
#### 4.5 Подписки
|
||||
- Подписка на пользователя: `AddBlock` с `ConnectionBody` подтип `CONNECTION_FOLLOW`, target = HEADER пользователя.
|
||||
- Подписка на канал: `AddBlock` с `ConnectionBody` подтип `CONNECTION_FOLLOW`, target = root блока канала (`CreateChannelBody` или HEADER для канала `0`).
|
||||
|
||||
### 5. API (форматы)
|
||||
|
||||
## 5.1 ListSubscriptionsFeed (чтение)
|
||||
Request:
|
||||
```json
|
||||
{
|
||||
"op": "ListSubscriptionsFeed",
|
||||
"requestId": "...",
|
||||
"payload": {
|
||||
"login": "A1",
|
||||
"limit": 200
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
Response (смысловые поля):
|
||||
- `ownedChannels[]`
|
||||
- `followedUsersChannels[]`
|
||||
- `followedChannels[]`
|
||||
|
||||
---
|
||||
|
||||
## 5.2 GetChannelMessages (чтение)
|
||||
Request:
|
||||
```json
|
||||
{
|
||||
"op": "GetChannelMessages",
|
||||
"requestId": "...",
|
||||
"payload": {
|
||||
"channel": {
|
||||
"ownerBlockchainName": "A1-001",
|
||||
"channelRootBlockNumber": 0,
|
||||
"channelRootBlockHash": ""
|
||||
},
|
||||
"limit": 200,
|
||||
"sort": "asc"
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
Response (смысловые поля):
|
||||
- `channel`
|
||||
- `messages[]`
|
||||
|
||||
---
|
||||
|
||||
## 5.3 GetMessageThread (чтение)
|
||||
Request:
|
||||
```json
|
||||
{
|
||||
"op": "GetMessageThread",
|
||||
"requestId": "...",
|
||||
"payload": {
|
||||
"message": {
|
||||
"blockchainName": "A1-001",
|
||||
"blockNumber": 15,
|
||||
"blockHash": "..."
|
||||
},
|
||||
"depthUp": 20,
|
||||
"depthDown": 2,
|
||||
"limitChildrenPerNode": 50
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
Response (смысловые поля):
|
||||
- `ancestors[]`
|
||||
- `focus`
|
||||
- `descendants[]`
|
||||
|
||||
---
|
||||
|
||||
## 5.4 AddBlock (запись)
|
||||
Любое изменение (создать канал, пост, reply, реакция, подписка) записывается через:
|
||||
```json
|
||||
{
|
||||
"op": "AddBlock",
|
||||
"requestId": "...",
|
||||
"payload": {
|
||||
"blockchainName": "A1-001",
|
||||
"blockNumber": 6,
|
||||
"prevBlockHash": "<64-hex>",
|
||||
"blockBytesB64": "<base64 full block>"
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
Важно:
|
||||
- `blockBytesB64` формируется на клиенте.
|
||||
- Подпись блока формируется на клиенте приватным blockchain key пользователя.
|
||||
- Перед добавлением блока клиент берет актуальный курсор цепочки с сервера.
|
||||
|
||||
### 6. Типы блоков для каналов и связей (через AddBlock)
|
||||
- Создание канала: `CreateChannelBody`
|
||||
- Пост/ответ: `TextBody` (`TEXT_POST`, `TEXT_REPLY`)
|
||||
- Реакции: `ReactionBody` (лайк/снятие лайка)
|
||||
- Подписки: `ConnectionBody` (`CONNECTION_FOLLOW`)
|
||||
|
||||
### 7. Критерии приемки
|
||||
- Список каналов отображается с реальными данными API.
|
||||
- Формат названия канала в UI: `ownerLogin/channelName`.
|
||||
- Создание канала реально пишет блок и канал появляется после обновления.
|
||||
- Отправка поста/ответа/реакций реально пишет блок и видна после перечитки API.
|
||||
- Подписка на пользователя/канал реально пишет блок и отражается в выдаче.
|
||||
- Переход в тред сообщения показывает реальные `ancestors/focus/descendants`.
|
||||
- Непрочитанные в списке каналов = общее число сообщений (временное правило).
|
||||
|
||||
### 8. Локальный запуск (уже сделано)
|
||||
Команда:
|
||||
```bash
|
||||
./gradlew startLocal
|
||||
```
|
||||
|
||||
Что делает:
|
||||
- чистит логи;
|
||||
- билдит сервер;
|
||||
- запускает локальный WS сервер;
|
||||
- запускает локальный HTTP сервер клиента;
|
||||
- открывает браузер по URL с параметром `localWsPort`.
|
||||
@@ -1,25 +0,0 @@
|
||||
# Краткое описание задачи
|
||||
|
||||
Нужно сделать полностью рабочую вкладку «Каналы» в SHiNE.
|
||||
|
||||
Пользователь должен:
|
||||
- видеть список каналов;
|
||||
- открывать канал и читать сообщения;
|
||||
- открывать тред сообщения;
|
||||
- отвечать, ставить и убирать лайк;
|
||||
- подписываться на пользователей и каналы.
|
||||
|
||||
Чтение данных идет через 3 API:
|
||||
- `ListSubscriptionsFeed`
|
||||
- `GetChannelMessages`
|
||||
- `GetMessageThread`
|
||||
|
||||
Все действия записи делаются только через `AddBlock` с подписью на клиенте.
|
||||
|
||||
Формат имени канала в интерфейсе:
|
||||
- `имя_пользователя/имя_канала`
|
||||
|
||||
Локальный запуск проекта:
|
||||
```bash
|
||||
./gradlew startLocal
|
||||
```
|
||||
@@ -1,153 +0,0 @@
|
||||
# Задача 02: Web Push + подписанный API отправки личных сообщений
|
||||
|
||||
## Контекст (по текущему состоянию проекта)
|
||||
- Уже есть JSON WebSocket API для личных сообщений: `SendDirectMessage`, `AckIncomingMessage`, `UpsertPushToken`.
|
||||
- Сейчас серверный fallback-пуш реализован через FCM (`FcmPushSender`) и ключ `fcm.server.key`.
|
||||
- Клиент уже регистрирует service worker и токен Firebase, затем аплоадит push token на сервер.
|
||||
|
||||
## Цель
|
||||
Добавить полностью рабочий сценарий доставки личных сообщений с приоритетом:
|
||||
1) онлайн-доставка в активную WebSocket-сессию;
|
||||
2) если не подтверждено — Web Push;
|
||||
3) поддержать отдельный API отправки без авторизации, где доступ проверяется цифровой подписью Ed25519 по `clientKey` отправителя.
|
||||
|
||||
---
|
||||
|
||||
## Предварительная спецификация подписанного пакета (v1)
|
||||
> ВАЖНО: финально фиксируется после уточнений по endian/кодировкам/лимитам.
|
||||
|
||||
Пакет (binary):
|
||||
1. `prefix` — ASCII-константа, например `SHINE_MESSAGE`.
|
||||
2. `toLoginLen` — 1 байт.
|
||||
3. `toLogin` — ASCII, длина = `toLoginLen`.
|
||||
4. `fromLoginLen` — 1 байт.
|
||||
5. `fromLogin` — ASCII, длина = `fromLoginLen`.
|
||||
6. `timeMs` — 8 байт (unix ms).
|
||||
7. `nonce32` — 4 байта случайное число.
|
||||
8. `messageType` — 4 байта.
|
||||
9. `targetMode` — 1 байт:
|
||||
- `0` = всем сессиям пользователя,
|
||||
- `1` = конкретной сессии.
|
||||
10. Если `targetMode=1`:
|
||||
- `sessionIdLen` — 1 байт,
|
||||
- `sessionId` — ASCII.
|
||||
11. `messageLen` — 2 байта.
|
||||
12. `messageBytes` — бинарные данные длиной `messageLen`.
|
||||
13. `signature64` — 64 байта, Ed25519 подпись всего блока **без** `signature64`.
|
||||
|
||||
Ограничения (первичный draft):
|
||||
- общий размер пакета ≤ 4000 байт;
|
||||
- логины/префикс/идентификатор сессии — ASCII;
|
||||
- повторы отсекаются по `(fromLogin, timeMs, nonce32)` в окне TTL.
|
||||
|
||||
---
|
||||
|
||||
## Сервер: что доработать
|
||||
|
||||
### 1) Новый endpoint без авторизации
|
||||
Операция (через WS JSON обертку) условно `SendSignedDirectMessage`:
|
||||
- принимает пакет (base64 binary blob);
|
||||
- парсит и валидирует формат;
|
||||
- достает `fromLogin`, поднимает `clientKey` пользователя;
|
||||
- проверяет подпись Ed25519;
|
||||
- проверяет анти-replay (time window + nonce);
|
||||
- отправляет сообщение по правилам маршрутизации;
|
||||
- пишет результат (messageId, каналы доставки, причины недоставки).
|
||||
|
||||
### 2) Маршрутизация доставки
|
||||
Для `targetMode=1`:
|
||||
- если целевая сессия онлайн и ACK пришел вовремя — успех;
|
||||
- иначе отправка в Web Push этой сессии (если есть subscription).
|
||||
|
||||
Для `targetMode=0`:
|
||||
- обход всех сессий пользователя;
|
||||
- сначала online delivery + ACK;
|
||||
- для непринятых/офлайн — Web Push по соответствующим subscription;
|
||||
- если subscription отсутствует — тихий skip.
|
||||
|
||||
### 3) Миграция от FCM к Web Push
|
||||
- добавить конфиг VAPID (`webpush.public.key`, `webpush.private.key`, `webpush.subject`);
|
||||
- хранить на сервере не только token, а web-push subscription (endpoint + keys);
|
||||
- сделать отправщик Web Push и заменить/расширить текущий `FcmPushSender`.
|
||||
|
||||
### 4) Безопасность
|
||||
- строгая ASCII-валидация логинов/sessionId;
|
||||
- лимиты длины всех полей;
|
||||
- rate limit на endpoint;
|
||||
- audit-лог неуспешных проверок подписи/формата;
|
||||
- защита от replay.
|
||||
|
||||
---
|
||||
|
||||
## Клиент (shine-UI): что доработать
|
||||
|
||||
1. Перейти на стандартный Web Push flow:
|
||||
- регистрация service worker;
|
||||
- `PushManager.subscribe(...)` с VAPID public key;
|
||||
- отправка subscription на сервер (`UpsertPushSubscription` или расширение `UpsertPushToken`).
|
||||
|
||||
2. Service worker:
|
||||
- `push` handler получает payload целиком;
|
||||
- показывает системное уведомление;
|
||||
- при клике открывает/фокусирует нужный чат.
|
||||
|
||||
3. Online-сообщения:
|
||||
- сохранить текущий event-канал `IncomingDirectMessage`;
|
||||
- обязателен ACK (`AckIncomingMessage` уже есть).
|
||||
|
||||
4. Keep-alive:
|
||||
- UI отправляет `Ping` раз в 60 секунд при активной сессии.
|
||||
|
||||
---
|
||||
|
||||
## Документация
|
||||
Сделать отдельный документ настройки Web Push:
|
||||
- как сгенерировать VAPID ключи;
|
||||
- какие параметры прописать на сервере и в UI;
|
||||
- как проверить локально e2e (онлайн + офлайн пуш);
|
||||
- ограничения payload и рекомендации по ретраям.
|
||||
|
||||
---
|
||||
|
||||
## Этапы реализации (предложение)
|
||||
1. Зафиксировать бинарный формат + валидации.
|
||||
2. Реализовать серверный parser/validator/signature verify/replay guard.
|
||||
3. Реализовать Web Push sender + storage subscription.
|
||||
4. Подключить новый endpoint и маршрутизацию доставки.
|
||||
5. Обновить UI (subscription + service worker + ping timer).
|
||||
6. Добавить интеграционные тесты (online ACK / offline push / bad signature / replay / oversize).
|
||||
7. Добавить документацию.
|
||||
|
||||
---
|
||||
|
||||
## Что нужно уточнить до разработки
|
||||
1. Endian для `timeMs/nonce/messageType/messageLen` (big-endian или little-endian).
|
||||
2. Что именно подписывается: строго весь префикс..messageBytes (без подписи) — подтвердить.
|
||||
3. Диапазон допустимых `messageType`.
|
||||
4. TTL окна для анти-replay (например 5 минут / 15 минут).
|
||||
5. Лимиты длин для login/session/message.
|
||||
6. Можно ли временно оставить FCM как fallback, пока не готов Web Push в проде.
|
||||
7. Формат сообщения в `messageBytes`: opaque bytes или UTF-8 строка.
|
||||
|
||||
|
||||
## Статус реализации (12.04.2026)
|
||||
|
||||
### Что уже внедрено в коде
|
||||
- `SendDirectMessage` переведён на signed-binary payload (`blobB64`) без обязательной авторизации WS-сессии.
|
||||
- Внедрён бинарный парсер пакета формата `SHiNE_msg + version(1) + ... + signature64`.
|
||||
- Проверка подписи Ed25519 делается по `clientKey` отправителя через `shine-server-crypto` (`Ed25519Util`).
|
||||
- Добавлен anti-replay guard `(from_login, time_ms, nonce)` с TTL 15 минут.
|
||||
- Добавлено историческое хранилище `signed_direct_messages_history` с сырым пакетом `raw_packet`.
|
||||
- Логика доставки: сначала WS+ACK, затем fallback на Web Push (по подписке конкретной session).
|
||||
- Поле типа сообщения переведено на `uint16`, пока поддерживается только `1`.
|
||||
- Для `targetMode=1` при несуществующей сессии возвращается `success` с `sessionNotFound=true` и `delivered=0`.
|
||||
- UI переведён с Firebase/FCM на браузерный `PushManager.subscribe` + Service Worker `push`.
|
||||
- Добавлен keep-alive ping из UI раз в 60 секунд при авторизованной сессии.
|
||||
|
||||
### Что настроить в окружении
|
||||
- В `application.properties` задать:
|
||||
- `webpush.vapid.public`
|
||||
- `webpush.vapid.private`
|
||||
- `webpush.vapid.subject`
|
||||
- В `shine-UI/index.html` задать публичный VAPID ключ в `window.__SHINE_WEBPUSH_VAPID_PUBLIC_KEY__`.
|
||||
|
||||
@@ -0,0 +1,86 @@
|
||||
# Third-party notices
|
||||
|
||||
This file lists third-party components and assets used by SHiNE. It is a practical attribution file for repository publication; individual dependency artifacts may include additional license text in their own packages.
|
||||
|
||||
## Twemoji
|
||||
|
||||
SHiNE UI uses Twemoji graphics for emoji rendering in the chat emoji picker and emoji-only messages.
|
||||
|
||||
- Source: https://github.com/twitter/twemoji
|
||||
- Graphics license: Creative Commons Attribution 4.0 International (CC BY 4.0)
|
||||
- Code license: MIT
|
||||
- License text: https://creativecommons.org/licenses/by/4.0/
|
||||
|
||||
Attribution:
|
||||
|
||||
Emoji graphics by Twemoji, licensed under CC BY 4.0.
|
||||
|
||||
Changes in SHiNE:
|
||||
|
||||
- Emoji are stored and sent as normal Unicode text.
|
||||
- Twemoji SVG graphics are used only as a visual rendering layer in the web UI.
|
||||
- The UI loads Twemoji SVG assets from the pinned `twitter/twemoji@14.0.2` package via jsDelivr.
|
||||
|
||||
## QR Code Generator
|
||||
|
||||
SHiNE UI includes `shine-UI/js/vendor-qrcode-generator.js`.
|
||||
|
||||
- Project: QR Code Generator for JavaScript
|
||||
- Copyright: Kazuhiko Arase
|
||||
- License: MIT
|
||||
- Source: http://www.d-project.com/
|
||||
|
||||
## SHiNE project-owned visual assets
|
||||
|
||||
SHiNE logos, icons and generated visual assets in `shine-UI/assets/` and `shine-UI/img/` are project-owned assets, unless a specific file says otherwise.
|
||||
|
||||
## Solana JavaScript libraries
|
||||
|
||||
SHiNE UI and Solana tooling use Solana JavaScript libraries, including `@solana/web3.js`.
|
||||
|
||||
- Package: `@solana/web3.js`
|
||||
- License: MIT
|
||||
- Source: https://github.com/solana-foundation/solana-web3.js
|
||||
|
||||
Some Solana JavaScript dependency trees include additional packages under permissive licenses such as MIT, Apache-2.0, BSD, ISC, 0BSD and CC0-1.0.
|
||||
|
||||
Known stricter dependency:
|
||||
|
||||
- Package: `rpc-websockets`
|
||||
- License: LGPL-3.0-only
|
||||
- Used transitively through Solana JavaScript dependencies in `shine-solana/shine` and `SHiNE-browser-plugin-wallet`.
|
||||
|
||||
## Noble cryptography libraries
|
||||
|
||||
SHiNE browser wallet/vendor bundles include Noble cryptography code.
|
||||
|
||||
- Packages: `@noble/curves`, `@noble/hashes`
|
||||
- License: MIT
|
||||
- Source: https://github.com/paulmillr/noble-curves and https://github.com/paulmillr/noble-hashes
|
||||
|
||||
## Java server dependencies
|
||||
|
||||
The SHiNE Java server uses third-party dependencies from Maven Central, including:
|
||||
|
||||
- Eclipse Jetty (`org.eclipse.jetty:*`) - EPL-2.0 / Apache-2.0 family licensing
|
||||
- Bouncy Castle (`org.bouncycastle:bcprov-jdk18on`) - Bouncy Castle permissive license
|
||||
- Jackson (`com.fasterxml.jackson.core:jackson-databind`) - Apache-2.0
|
||||
- Logback (`ch.qos.logback:logback-classic`) - EPL-1.0 / LGPL-2.1
|
||||
- SLF4J (`org.slf4j:slf4j-api`) - MIT
|
||||
- SQLite JDBC (`org.xerial:sqlite-jdbc`) - Apache-2.0
|
||||
- Web Push Java library (`nl.martijndwars:web-push`) - Apache-2.0
|
||||
- JUnit (`org.junit:*`) - EPL-2.0
|
||||
|
||||
## Gradle Wrapper
|
||||
|
||||
The repository includes Gradle wrapper scripts.
|
||||
|
||||
- License: Apache License 2.0
|
||||
- Source: https://gradle.org/
|
||||
|
||||
## Espressif code snippets
|
||||
|
||||
The ESP32 prototype area includes Espressif audio codec helper files with SPDX/license headers.
|
||||
|
||||
- Files include `ESPRESSIF MIT License` and `SPDX-License-Identifier: Apache-2.0` notices.
|
||||
- Original copyright notices are preserved in the source files.
|
||||
@@ -35,6 +35,5 @@
|
||||
|
||||
## Какие документы потом обновить
|
||||
|
||||
- `Dev_Docs/deploy/`;
|
||||
- `Dev_Docs/Blockchain/sync-between-servers.md`, если изменится поведение остановки/восстановления.
|
||||
|
||||
- `deploy/`;
|
||||
- `docs/Blockchain/sync-between-servers.md`, если изменится поведение остановки/восстановления.
|
||||
|
||||
@@ -24,6 +24,6 @@ QR-подключение других устройств сейчас есть
|
||||
|
||||
## Какие документы потом обновить
|
||||
|
||||
- `Dev_Docs/Solana_Architecture/README.md`;
|
||||
- `docs/Solana_Architecture/README.md`;
|
||||
- `TODO/medium/2026-06-03_подключение_других_устройств_через_qr.md`;
|
||||
- `TODO/medium/2026-06-02_сессионные_homeserver_в_pda.md`.
|
||||
|
||||
@@ -12,16 +12,30 @@
|
||||
- откуда продолжать;
|
||||
- какие документы потом надо обновить.
|
||||
- Это не активная разработка. Тут только план и контекст.
|
||||
- Старую папку `Dev_Docs/Future_Features/` считать архивной и больше не использовать как источник новых задач.
|
||||
- Старую папку `docs/Future_Features/` считать архивной и больше не использовать как источник новых задач.
|
||||
|
||||
## Текущие задачи
|
||||
|
||||
- `2026-06-26_1800_корректное_завершение_за_30с.md` - дать сервису до 30 секунд на корректное завершение опасных операций перед рестартом.
|
||||
- `2026-06-26_1805_межсерверный_ws_и_dm_sync.md` - постоянный server-to-server WebSocket, push новых блоков и DM, ACK и backfill.
|
||||
- `2026-06-26_1810_подключение_устройств_по_qr.md` - довести подключение других устройств по QR и перевести это в нормальные типизированные сессии.
|
||||
- `2026-06-26_1815_esp32_файловое_хранилище.md` - использовать ESP32 как личное файловое хранилище для переписок и вложений.
|
||||
|
||||
## Перенесённые планы из `Dev_Docs/Future_Features/`
|
||||
## Децентрализация
|
||||
|
||||
Текущий production-режим SHiNE считается односерверным. Задачи по нескольким серверам, Arweave и realtime PDA/Solana sync вынесены в `Децентрализация/` и не блокируют выкладку текущей версии на GitHub.
|
||||
|
||||
- `Децентрализация/односерверный_production_режим.md` - границы текущей production-версии с одним сервером.
|
||||
- `Децентрализация/запись_блокчейнов_в_arweave.md` - будущая запись/архивация блокчейнов в Arweave.
|
||||
- `Децентрализация/realtime_pda_solana_sync.md` - будущая онлайн-синхронизация PDA и Solana.
|
||||
- `Децентрализация/межсерверная_передача_сообщений.md` - будущая доставка сообщений между серверами.
|
||||
- `Децентрализация/межсерверные_звонки.md` - будущая маршрутизация звонков между серверами.
|
||||
- `Децентрализация/2026-06-26_1805_межсерверный_ws_и_dm_sync.md` - перенесённый старый план постоянного server-to-server WS и DM sync.
|
||||
|
||||
## Новые фишки которые надо доделать
|
||||
|
||||
- `Новые фишки которые надо доделать/Новая_контентная_модель_блокчейна/` - отложенная новая контентная модель блокчейна, не входящая в текущий односерверный production-релиз.
|
||||
|
||||
## Перенесённые планы из `docs/Future_Features/`
|
||||
|
||||
### near
|
||||
|
||||
@@ -35,6 +49,7 @@
|
||||
- `medium/2026-05-26_0029_esp32s3_file_storage.md` - ESP32S3 как личное файловое хранилище SHiNE для файлов переписок и вложений.
|
||||
- `medium/2026-06-02_сессионные_homeserver_в_pda.md` - несколько homeserver-ов пользователя как типизированные сессии в PDA с версией записи.
|
||||
- `medium/2026-06-03_подключение_других_устройств_через_qr.md` - довести подключение других устройств через QR: сейчас заготовка есть, но сценарий работает нестабильно и его нужно будет отдельно доделать.
|
||||
- `medium/2026-07-22_переход_с_sqlite_на_postgresql.md` - подготовить перевод серверной БД с `SQLite` на `PostgreSQL` для более серьёзной конкурентной нагрузки и дальнейшего масштабирования.
|
||||
|
||||
### dao_запуск
|
||||
|
||||
|
||||
@@ -1,64 +0,0 @@
|
||||
# ESP32 как аппаратный кошелёк (device-сессия)
|
||||
|
||||
## Суть фичи
|
||||
|
||||
ESP32 становится аппаратным HSM (hardware security module): хранит ключи, постоянно подключён к SHiNE-серверу как device-сессия, подтверждает операции нажатием на экране. Другие устройства (браузер, телефон) взаимодействуют с ESP32 через сервер — без прямого соединения.
|
||||
|
||||
## Два ключевых сценария
|
||||
|
||||
### Сценарий 1 — Создание делегированной сессии
|
||||
1. Браузер/телефон → сервер: «хочу делегированную сессию от имени пользователя X»
|
||||
2. Сервер → ESP32 (device-сессия): «запрос на одобрение»
|
||||
3. Пользователь нажимает «Да» на сенсорном экране ESP32
|
||||
4. ESP32 → сервер: одобрено → сервер создаёт делегированную сессию для браузера
|
||||
|
||||
### Сценарий 2 — Подпись транзакции / блока
|
||||
1. Браузер (через делегированную сессию) → сервер → ESP32: «подпиши вот это»
|
||||
2. ESP32 показывает запрос на экране, пользователь подтверждает
|
||||
3. ESP32 подписывает нужным ключом → ответ через сервер → браузер
|
||||
|
||||
## Что нужно сделать
|
||||
|
||||
### ESP32 (основная работа)
|
||||
- [ ] Инициализация WiFi (SSID/пароль в NVS)
|
||||
- [ ] WebSocket-клиент (`WebSocketsClient`) — постоянное соединение с сервером
|
||||
- [ ] Авторизация на сервере: `AuthChallenge` → `CreateAuthSession` через `clientKey` (уже есть в NVS), сохранить `sessionId` в NVS
|
||||
- [ ] Обработчик входящих WebSocket-событий: JSON-парсинг, диспетчер по типу
|
||||
- [ ] Новые UI-экраны: «Разрешить сессию?» и «Подписать?» с кнопками Да/Нет
|
||||
- [ ] Расширенное хранилище ключей в NVS (произвольные именованные ключи сверх базовых трёх)
|
||||
- [ ] Переподключение при разрыве (reconnect loop)
|
||||
|
||||
### Сервер (минимальные изменения)
|
||||
- [ ] Добавить поле `sessionType` (`USER` / `DEVICE`) в таблицу `active_sessions`
|
||||
- [ ] Новая операция `DeviceApprovalRequest` — браузер запрашивает одобрение у device-сессии
|
||||
- [ ] Новая операция `DeviceApprovalResponse` — ESP32 отвечает (одобрено/отклонено)
|
||||
- [ ] Новые операции `SignRequest` / `SignResponse` — запрос подписи и ответ
|
||||
- [ ] Роутинг: при получении запроса найти device-сессию через `ActiveConnectionsRegistry.getByLogin(login)` + фильтр по `sessionType=DEVICE`, переслать туда
|
||||
|
||||
### Клиент (отдельный этап)
|
||||
- [ ] Браузерное расширение или UI: создание делегированной сессии, отправка `SignRequest`
|
||||
|
||||
## Что уже готово (переиспользуем)
|
||||
|
||||
- **Роутинг сообщений** — `SendDirectMessage` с `TARGET_ONE_SESSION` и `CallSignalToSession` уже умеют точечно доставлять в конкретный `sessionId`. Механизм готов, нужно добавить только новые op-коды поверх него.
|
||||
- **Ed25519 на ESP32** — библиотека `<Ed25519.h>` уже используется в скетче. Подписи работают.
|
||||
- **NVS** — уже хранит логин, мастер-секрет, 3 пары ключей. Расширяется легко.
|
||||
- **`ActiveConnectionsRegistry`** — поиск по `login` и `sessionId` уже есть на сервере.
|
||||
- **Аутентификация** — схема `AuthChallenge` → `CreateAuthSession` через Ed25519 уже полностью реализована.
|
||||
|
||||
## Оценка сложности
|
||||
|
||||
| Компонент | Сложность |
|
||||
|---|---|
|
||||
| ESP32: WiFi + WebSocket-клиент + авторизация | Средняя |
|
||||
| ESP32: обработчик входящих + UI подтверждений | Средняя |
|
||||
| Сервер: флаг sessionType + 4 новых op-а + роутинг | Низкая–средняя |
|
||||
| Браузерное расширение | Высокая (отдельный этап) |
|
||||
|
||||
**Итого фазы ESP32 + сервер: ~1–1.5 недели.**
|
||||
|
||||
## С чего начинать
|
||||
|
||||
1. Серверная часть проще и быстрее — начать с добавления `sessionType` и `DeviceApprovalRequest/Response`.
|
||||
2. Затем ESP32: WiFi → WebSocket → авторизация → обработчик входящих → UI.
|
||||
3. Браузерное расширение — отдельная итерация после того как ESP32 + сервер работают.
|
||||
@@ -104,9 +104,9 @@
|
||||
|
||||
## Какие документы нужно будет обновить при реализации
|
||||
|
||||
- `Dev_Docs/Blockchain/README.md` и связанные файлы, если изменятся типы служебных сообщений или форматы блокчейн-команд.
|
||||
- `Dev_Docs/API/` если изменится публичный серверный API или появятся новые операции.
|
||||
- `Dev_Docs/Personal_Messages/Протокол_DM_v1.md` если часть маршрутизации или подтверждений будет встроена в существующую логику доставки/сессий.
|
||||
- `docs/Blockchain/README.md` и связанные файлы, если изменятся типы служебных сообщений или форматы блокчейн-команд.
|
||||
- `docs/API/` если изменится публичный серверный API или появятся новые операции.
|
||||
- `docs/Personal_Messages/Протокол_DM_v1.md` если часть маршрутизации или подтверждений будет встроена в существующую логику доставки/сессий.
|
||||
- Документацию по homeserver/ESP32, если появится пользовательская или сервисная файловая логика на устройстве.
|
||||
|
||||
## С какого места продолжать позже
|
||||
|
||||
@@ -56,11 +56,9 @@
|
||||
- Код формирования репоста в `auth-service.js` не удалён: его можно будет использовать как основу при возвращении к задаче.
|
||||
- Код отображения target-полей и перехода к оригиналу не удалён: он нужен для будущей проверки и возможной совместимости с уже созданными тестовыми блоками.
|
||||
|
||||
## Почему это не лежит в Pending_Features
|
||||
## Почему это лежит в TODO
|
||||
|
||||
`Dev_Docs/Pending_Features/` предназначена для фич, которые уже реализованы и ждут ручной проверки.
|
||||
|
||||
Репосты сейчас не подходят под этот статус: они не должны проверяться как готовая фича, потому что пользовательский сценарий временно закрыт, а серверная запись новых репостов заблокирована. Поэтому старый pending-файл удалён, а задача перенесена сюда как будущая.
|
||||
Репосты сейчас не должны проверяться как готовая фича, потому что пользовательский сценарий временно закрыт, а серверная запись новых репостов заблокирована. Поэтому задача остаётся в TODO как будущая.
|
||||
|
||||
## Что сделать при возврате к реализации
|
||||
|
||||
@@ -82,11 +80,11 @@
|
||||
- отображение `targetBlockchainName`, `targetBlockNumber`, `targetBlockHash`.
|
||||
7. Добавить или обновить тесты на успешный репост и отказ некорректных target-полей.
|
||||
8. Обновить документацию:
|
||||
- `Dev_Docs/Blockchain/11_TEXT_Blocks.md`;
|
||||
- `Dev_Docs/Blockchain/CHANGELOG.md`;
|
||||
- `Dev_Docs/API/04_Add_Block_to_Blockchain_API.md`;
|
||||
- `docs/Blockchain/11_TEXT_Blocks.md`;
|
||||
- `docs/Blockchain/CHANGELOG.md`;
|
||||
- `docs/API/04_Add_Block_to_Blockchain_API.md`;
|
||||
- документы API чтения каналов/тредов, если изменятся поля ответа.
|
||||
9. После реализации перенести задачу из `TODO/` в `Dev_Docs/Pending_Features/` как фичу, требующую ручной проверки.
|
||||
9. После реализации отдельно согласовать ручную проверку пользовательского сценария.
|
||||
|
||||
## Минимальный чек-лист ручной проверки в будущем
|
||||
|
||||
|
||||
@@ -47,10 +47,10 @@
|
||||
|
||||
## Документы, которые обновить при реализации
|
||||
|
||||
- `Dev_Docs/Blockchain/`, если появятся или изменятся блоки баланса.
|
||||
- `Dev_Docs/Blockchain/CHANGELOG.md`, если меняется блокчейн-формат.
|
||||
- `Dev_Docs/API/`, если меняется серверный API.
|
||||
- `Dev_Docs/Pending_Features/` - добавить файл ручной проверки после реализации.
|
||||
- `docs/Blockchain/`, если появятся или изменятся блоки баланса.
|
||||
- `docs/Blockchain/CHANGELOG.md`, если меняется блокчейн-формат.
|
||||
- `docs/API/`, если меняется серверный API.
|
||||
- после реализации отдельно согласовать ручную проверку.
|
||||
- Документацию Solana-регистрации, если баланс будет связан с Solana-модулем.
|
||||
|
||||
## Минимальная проверка в будущем
|
||||
|
||||
@@ -34,10 +34,10 @@
|
||||
|
||||
## Документы, которые нужно обновить при возврате
|
||||
|
||||
- `Dev_Docs/Keys/README.md`
|
||||
- `Dev_Docs/Personal_Messages/Протокол_DM_v1.md`
|
||||
- `Dev_Docs/API/`
|
||||
- `Dev_Docs/Blockchain/`, если появятся новые блоки или команды для файлов.
|
||||
- `docs/Keys/README.md`
|
||||
- `docs/Personal_Messages/Протокол_DM_v1.md`
|
||||
- `docs/API/`
|
||||
- `docs/Blockchain/`, если появятся новые блоки или команды для файлов.
|
||||
|
||||
## С какого места продолжать
|
||||
|
||||
|
||||
@@ -84,11 +84,11 @@
|
||||
## Что нужно обновить при реализации
|
||||
|
||||
- `shine-solana/shine/doc/formats/shine-user-pda-format-v.1.0.md`
|
||||
- `Dev_Docs/Solana_Architecture/README.md`
|
||||
- `Dev_Docs/Инициализация_Solana_регистрации/README.md`
|
||||
- `Dev_Docs/Keys/README.md`
|
||||
- `Dev_Docs/Personal_Messages/Протокол_DM_v1.md`, если изменится адресация DM по типам сессий
|
||||
- `Dev_Docs/API/`, если появятся новые серверные операции или изменятся ответы
|
||||
- `docs/Solana_Architecture/README.md`
|
||||
- `docs/Инициализация_Solana_регистрации/README.md`
|
||||
- `docs/Keys/README.md`
|
||||
- `docs/Personal_Messages/Протокол_DM_v1.md`, если изменится адресация DM по типам сессий
|
||||
- `docs/API/`, если появятся новые серверные операции или изменятся ответы
|
||||
|
||||
## Что пока не делать
|
||||
|
||||
|
||||
@@ -37,9 +37,8 @@
|
||||
|
||||
## Что обновить при возврате
|
||||
|
||||
- `Dev_Docs/Pending_Features/README.md`
|
||||
- после реализации отдельно согласовать ручную проверку
|
||||
- `shine-UI/js/pages/connect-device-view.js`
|
||||
- `shine-UI/js/pages/device-qr-view.js`
|
||||
- `shine-UI/js/services/qr-key-transfer-service.js`
|
||||
- документацию по ключам, если формат переноса меняется
|
||||
|
||||
|
||||
@@ -0,0 +1,36 @@
|
||||
# Переход с SQLite на PostgreSQL
|
||||
|
||||
## Зачем
|
||||
|
||||
Текущая серверная база на `SQLite` удобна для простого односерверного режима, но она хуже подходит для большого числа параллельных записей, роста нагрузки и дальнейшего масштабирования сервера.
|
||||
|
||||
`PostgreSQL` нужен как следующий уровень серверной БД для более надёжной конкурентной записи, более предсказуемой работы под нагрузкой и дальнейшего роста проекта.
|
||||
|
||||
## Что сделать
|
||||
|
||||
- Подготовить план переноса серверной БД с `SQLite` на `PostgreSQL`.
|
||||
- Найти все места, где код завязан на особенности `SQLite`.
|
||||
- Проверить все DAO и SQL-запросы на совместимость с `PostgreSQL`.
|
||||
- Продумать схему миграции существующей production/test базы без потери данных.
|
||||
- Отдельно проверить транзакции, `UPSERT`, индексы, case-insensitive сравнения и миграции схемы.
|
||||
- После этого подготовить отдельный этап внедрения и переключения сервера.
|
||||
|
||||
## Что уже есть в коде
|
||||
|
||||
- Доступ к БД в основном проходит через DAO-слой, а не полностью размазан по проекту.
|
||||
- Основная серверная логика уже разделена по модулям.
|
||||
- Но SQL и миграции сейчас написаны под `SQLite` и потребуют отдельного прохода.
|
||||
|
||||
## Откуда продолжать
|
||||
|
||||
- Начать с инвентаризации всех DAO и схемы БД.
|
||||
- После этого сделать отдельный документ с оценкой объёма работ по переносу.
|
||||
- Затем решить, будет ли это:
|
||||
- полный перевод сервера на `PostgreSQL`;
|
||||
- или поддержка двух драйверов на переходный период.
|
||||
|
||||
## Что потом обновить
|
||||
|
||||
- Серверную документацию по БД и миграциям.
|
||||
- Инструкции по локальному запуску сервера.
|
||||
- Скрипты деплоя и настройки окружения.
|
||||
@@ -58,8 +58,8 @@
|
||||
## Документы, которые обновить при реализации
|
||||
|
||||
- Документацию UI/кошельков, если такая есть.
|
||||
- `Dev_Docs/Pending_Features/` - добавить файл ручной проверки после реализации.
|
||||
- `Dev_Docs/API/`, только если появится новый серверный API или логирование.
|
||||
- после реализации отдельно согласовать ручную проверку.
|
||||
- `docs/API/`, только если появится новый серверный API или логирование.
|
||||
|
||||
## Минимальная проверка
|
||||
|
||||
|
||||
@@ -69,7 +69,7 @@ git show 0240db5:shine-solana/shine/programs/shine_login_guard/src/lib.rs
|
||||
- `shine-solana/shine/doc/programs/shine_login_guard.md`
|
||||
|
||||
5. Архитектурная документация:
|
||||
- `Dev_Docs/Solana_Architecture/README.md`
|
||||
- `docs/Solana_Architecture/README.md`
|
||||
|
||||
6. UI-логика precheck:
|
||||
- `shine-UI/js/pages/register-view.js`
|
||||
@@ -108,7 +108,7 @@ git show 0240db5:shine-solana/shine/programs/shine_login_guard/src/lib.rs
|
||||
4. Проверить, что `build.rs` всё ещё генерирует `generated_dictionary.rs` в прежнем формате.
|
||||
5. Сверить актуальность словарей в `src/dictionaries`.
|
||||
6. Обновить `shine_login_guard.md` обратно под словарную логику.
|
||||
7. Обновить `Dev_Docs/Solana_Architecture/README.md`.
|
||||
7. Обновить `docs/Solana_Architecture/README.md`.
|
||||
|
||||
### Проверка после возврата
|
||||
|
||||
|
||||
@@ -28,6 +28,6 @@
|
||||
|
||||
## Какие документы потом обновить
|
||||
|
||||
- `Dev_Docs/Personal_Messages/Протокол_DM_v1.md`, если изменится стратегия миграции старой истории;
|
||||
- `Dev_Docs/Personal_Messages/Формат_DM_v1.md`, если появится отдельное правило совместимости/конвертации;
|
||||
- при необходимости `Dev_Docs/API/12_Direct_Messages_Push_Calls_API.md`, если затронется поведение backlog/доставки.
|
||||
- `docs/Personal_Messages/Протокол_DM_v1.md`, если изменится стратегия миграции старой истории;
|
||||
- `docs/Personal_Messages/Формат_DM_v1.md`, если появится отдельное правило совместимости/конвертации;
|
||||
- при необходимости `docs/API/12_Direct_Messages_Push_Calls_API.md`, если затронется поведение backlog/доставки.
|
||||
|
||||
@@ -2,7 +2,9 @@
|
||||
|
||||
## Зачем
|
||||
|
||||
Сейчас синхронизация между серверами работает в основном как periodic sync и one-shot push. Для нормальной репликации ещё нужен постоянный межсерверный канал:
|
||||
Текущий production-режим SHiNE рассчитан на один основной сервер. Межсерверная синхронизация относится к будущей децентрализации и не должна блокировать выкладку односерверной production-версии.
|
||||
|
||||
Сейчас синхронизация между серверами работает в основном как periodic sync и one-shot push. Для нормальной репликации в будущем ещё нужен постоянный межсерверный канал:
|
||||
|
||||
- живое подключение к партнёру;
|
||||
- push новых блоков;
|
||||
@@ -35,7 +37,7 @@
|
||||
|
||||
## Какие документы потом обновить
|
||||
|
||||
- `Dev_Docs/Blockchain/sync-between-servers.md`;
|
||||
- `Dev_Docs/Personal_Messages/Протокол_DM_v1.md`;
|
||||
- `Dev_Docs/Personal_Messages/Формат_DM_v1.md`;
|
||||
- `Dev_Docs/API/`.
|
||||
- `docs/Blockchain/sync-between-servers.md`;
|
||||
- `docs/Personal_Messages/Протокол_DM_v1.md`;
|
||||
- `docs/Personal_Messages/Формат_DM_v1.md`;
|
||||
- `docs/API/`.
|
||||
@@ -0,0 +1,16 @@
|
||||
# Децентрализация
|
||||
|
||||
Папка для задач, которые нужны для будущего режима с несколькими серверами, Solana/PDA-синхронизацией и внешним хранением данных.
|
||||
|
||||
## Текущий статус
|
||||
|
||||
Сейчас production-режим SHiNE считается односерверным: один сервер обслуживает пользователей, сообщения, звонки и запись данных. Задачи из этой папки не являются блокерами для выкладки текущего репозитория на GitHub и запуска одного production-сервера.
|
||||
|
||||
## Задачи
|
||||
|
||||
- `односерверный_production_режим.md` - зафиксировать границы текущей production-версии.
|
||||
- `запись_блокчейнов_в_arweave.md` - вынести долговременную запись блокчейнов в Arweave.
|
||||
- `realtime_pda_solana_sync.md` - сделать онлайн-синхронизацию PDA/Solana в реальном времени.
|
||||
- `межсерверная_передача_сообщений.md` - реализовать доставку сообщений между серверами.
|
||||
- `межсерверные_звонки.md` - реализовать маршрутизацию звонков между серверами.
|
||||
- `2026-06-26_1805_межсерверный_ws_и_dm_sync.md` - старый план постоянного server-to-server WS и DM sync, перенесённый в контекст децентрализации.
|
||||
@@ -0,0 +1,30 @@
|
||||
# Realtime-синхронизация PDA и Solana
|
||||
|
||||
## Зачем
|
||||
|
||||
В будущем PDA-записи и Solana-состояние должны автоматически и быстро синхронизироваться с серверным состоянием, чтобы данные пользователей, homeserver-сессии и связанные записи не расходились.
|
||||
|
||||
## Что сделать
|
||||
|
||||
1. Определить, какие серверные события должны обновлять PDA.
|
||||
2. Добавить очередь/воркер для надёжной отправки изменений в Solana.
|
||||
3. Добавить периодическую сверку серверного состояния с PDA.
|
||||
4. Добавить обработку ошибок, повторов и конфликтов версий.
|
||||
5. Добавить мониторинг задержек и неуспешных Solana-транзакций.
|
||||
|
||||
## Что учесть
|
||||
|
||||
- Solana/Anchor-модуль находится в `shine-solana/shine/` и ведётся отдельно от основного server/UI deploy.
|
||||
- Перед изменениями внутри Solana-модуля нужно читать `shine-solana/shine/AGENTS.md`.
|
||||
- Основная инструкция по Solana-регистрации находится в `docs/Инициализация_Solana_регистрации/README.md`.
|
||||
- Формат пользовательской PDA-записи описан в `shine-solana/shine/doc/formats/shine-user-pda-format-v.1.0.md`.
|
||||
|
||||
## Документы, которые потом нужно обновить
|
||||
|
||||
- `docs/Инициализация_Solana_регистрации/README.md`;
|
||||
- `docs/Solana_Architecture/README.md`;
|
||||
- `shine-solana/shine/doc/formats/shine-user-pda-format-v.1.0.md`, если меняется формат PDA.
|
||||
|
||||
## Статус
|
||||
|
||||
Отложено до этапа децентрализации.
|
||||
@@ -0,0 +1,29 @@
|
||||
# Запись блокчейнов в Arweave
|
||||
|
||||
## Зачем
|
||||
|
||||
Для будущей децентрализации нужно долговременное внешнее хранение блокчейнов, чтобы данные не зависели только от одного серверного диска.
|
||||
|
||||
## Что сделать
|
||||
|
||||
1. Определить, какие блокчейны и какие диапазоны блоков записываются в Arweave.
|
||||
2. Зафиксировать формат пачки блоков, метаданных, ссылок и контрольных хэшей.
|
||||
3. Добавить безопасный механизм публикации без хранения приватного JWK в git.
|
||||
4. Добавить проверку уже загруженных диапазонов, чтобы не плодить дубли.
|
||||
5. Описать восстановление блокчейна из Arweave при потере локальных данных.
|
||||
|
||||
## Важные ограничения
|
||||
|
||||
- Любое изменение формата блокчейна требует отдельного предупреждения и явного подтверждения пользователя.
|
||||
- Добавление данных в блокчейн должно выполняться только через `AddBlock`.
|
||||
- Секреты Arweave нельзя хранить в репозитории.
|
||||
|
||||
## Документы, которые потом нужно обновить
|
||||
|
||||
- `docs/Blockchain/README.md`;
|
||||
- `docs/Blockchain/CHANGELOG.md`;
|
||||
- документы deploy/секретов в `deploy/`, если появятся новые параметры.
|
||||
|
||||
## Статус
|
||||
|
||||
Отложено до этапа децентрализации.
|
||||
@@ -0,0 +1,30 @@
|
||||
# Межсерверная передача сообщений
|
||||
|
||||
## Зачем
|
||||
|
||||
Когда у SHiNE появится несколько серверов, пользователи на разных серверах должны получать личные сообщения без ручной синхронизации и без привязки к одному центральному узлу.
|
||||
|
||||
## Что сделать
|
||||
|
||||
1. Определить протокол server-to-server доставки DM.
|
||||
2. Добавить маршрутизацию получателя по серверу, user id, публичному ключу или PDA.
|
||||
3. Добавить ACK, повторы, дедупликацию и backfill пропущенных сообщений.
|
||||
4. Разделить realtime-доставку и восстановление истории.
|
||||
5. Описать поведение при недоступности удалённого сервера.
|
||||
|
||||
## Что учесть
|
||||
|
||||
- Логика DM должна соответствовать документам в `docs/Personal_Messages/`.
|
||||
- При изменении формата signed DM-блока или правил доставки нужно обновлять протокол и байтовый формат DM.
|
||||
- Если появятся новые server API/WebSocket операции, нужно обновить `docs/API/`.
|
||||
|
||||
## Документы, которые потом нужно обновить
|
||||
|
||||
- `docs/Personal_Messages/Протокол_DM_v1.md`;
|
||||
- `docs/Personal_Messages/Формат_DM_v1.md`;
|
||||
- `docs/API/`;
|
||||
- `docs/API/09_Operations_Index.md`, если добавляются новые `op`.
|
||||
|
||||
## Статус
|
||||
|
||||
Отложено до этапа децентрализации.
|
||||
@@ -0,0 +1,29 @@
|
||||
# Межсерверные звонки
|
||||
|
||||
## Зачем
|
||||
|
||||
В будущем пользователи на разных серверах должны иметь возможность устанавливать звонки так же, как пользователи одного сервера.
|
||||
|
||||
## Что сделать
|
||||
|
||||
1. Определить протокол межсерверной сигнализации звонков.
|
||||
2. Добавить маршрутизацию offer/answer/ICE-кандидатов между серверами.
|
||||
3. Добавить обработку статусов занятости, отказа, таймаута и ошибок маршрута.
|
||||
4. Добавить диагностику доставки сигналов между серверами.
|
||||
5. Проверить совместимость с текущими логами `CallDeliveryReport`.
|
||||
|
||||
## Что учесть
|
||||
|
||||
- Специальная диагностика установки звонков идёт через `CallDeliveryReport`.
|
||||
- На production важно сохранять поля `reason`, `failureStage`, `pcConnectionState`, `pcIceConnectionState`, `routeLabel`, `configuredTurnHosts*`, `reachableTurnHosts*`.
|
||||
- Межсерверные звонки не должны ломать текущий односерверный сценарий.
|
||||
|
||||
## Документы, которые потом нужно обновить
|
||||
|
||||
- `docs/API/`, если добавляются или меняются операции сигнализации;
|
||||
- документы по звонкам/диагностике, если они будут выделены отдельно;
|
||||
- deploy-документы, если появятся новые параметры TURN/server-to-server маршрутизации.
|
||||
|
||||
## Статус
|
||||
|
||||
Отложено до этапа децентрализации.
|
||||