🔌 Доставка по gRPC
Decoded Shred Stream можно потреблять по gRPC вместо UDP: исходящий, ориентированный на соединение стрим, который добавляет фильтрацию по аккаунтам на стороне сервера и упорядоченную доставку по HTTP/2 поверх тех же транзакций с субмиллисекундной задержкой — без необходимости в публично доступной конечной точке UDP. В дашборде показаны gRPC-эндпоинт и токен доступа для каждого стрима.
- Эндпоинт —
<host>:50051(один порт для всех потоков). - Транспорт — gRPC поверх HTTP/2, в открытом виде h2c (без TLS). Это осознанный выбор ради задержки: не нужно платить за рукопожатие TLS. Изолируйте канал на сетевом уровне — приватная сеть, VPN или пиринг — и никогда не открывайте его в публичный интернет.
- Маршруты —
shredstream.com.DecodedShredStreamService/SubscribeDecodedTransactionsили совместимый «из коробки»shreder_binary.ShrederBinaryService/SubscribeBinaryTransactions. Оба строго эквивалентны (те же protobuf-сообщения, побайтово идентичное содержимое) — клиенту, приходящему с эндпоинта Shreder / Raiden Pulse, достаточно сменить адрес.
🔑 Аутентификация
Каждый вызов должен передавать ваш токен в метаданных gRPC в одной из двух форм (принимаются обе):
authorization: Bearer <TOKEN> x-token: <TOKEN>
Отсутствующий или недействительный токен отклоняется единым статусом UNAUTHENTICATED с сообщением authentication refused. Проверьте свой токен.
Помните: одно соединение на токен. Открытие второго стрима с тем же токеном вытесняет самый старый.
🧭 Модель подписки
Маршрут — это двунаправленный стриминг (stream request → stream response):
- Откройте стрим и отправьте хотя бы один запрос, содержащий map именованных фильтров:
{ "<name>": <filter>, … }. - Сервер доставляет только те транзакции, что совпадают хотя бы с одним фильтром. Каждый ответ помечается именем(-ами) сработавших фильтров (поле
filters) — так вы можете маршрутизировать несколько стратегий по одному стриму. - Вы можете отправить новую map в любой момент: она заменяет предыдущую на лету, без переподключения и без разрыва в потоке.
Крайние случаи:
- Пустая map (
{}) ⇒ ничего не доставляется (вы подписываетесь на фильтры, а не на firehose). - Пустой именованный фильтр (
{ "all": {} }) ⇒ проходит каждая транзакция, с меткой"all". - После отправки фильтров отправлять больше нечего: стрим продолжает доставку с этими фильтрами, пока вы не отправите новые или не закроете соединение.
🎯 Фильтры по аккаунтам
Каждый именованный фильтр — это три списка публичных ключей в base58, объединённых логическим И:
| Поле | Значение |
|---|---|
account_include | транзакция должна затрагивать хотя бы один из этих аккаунтов (пустой = без ограничений) |
account_exclude | транзакция не должна затрагивать ни один из этих аккаунтов |
account_required | транзакция должна затрагивать все эти аккаунты |
Для транзакций «затрагивание» оценивается по статическим ключам аккаунтов транзакции (включая подписантов). Адреса, разрешаемые через Address Lookup Tables (ALT), отсутствуют в конверте транзакции и потому не поддаются фильтрации — фильтруйте только по статическим ключам.
📦 Полезная нагрузка
Каждый ответ несёт транзакцию в виде сырых байтов вместе с её слотом:
message SubscribeUpdateBinaryTransaction {BinaryTransaction transaction = 1;uint64 slot = 2; // слот Solana}message BinaryTransaction {repeated bytes signatures = 1; // подписи (по 64 байта каждая)bytes binary_transaction = 3; // VersionedTransaction, bincode, VERBATIM}
binary_transaction — это стандартный wire-формат Solana — сериализованная через bincode VersionedTransaction, без изменений. Десериализуйте её любым SDK Solana и передайте вашему существующему парсеру, ровно как в режиме UDP. Транзакции голосования по умолчанию исключены.
💻 Официальный клиент
Официальные клиенты decoded-shredstream берут этот маршрут на себя — соединение, токен в метаданных, map фильтров, переподключение с повторной отправкой текущих фильтров и транзакцию вместе с её слотом и подписями. Один и тот же пакет для UDP и gRPC.
// npm install decoded-shredstreamimport { DecodedShredStream, FilterAll } from "decoded-shredstream";const client = await DecodedShredStream.grpc({endpoint: "<host>:50051",token: process.env.DECODED_SHREDSTREAM_TOKEN!,filters: { all: FilterAll }, // or: { "watched-wallet": { include: [wallet] } }});for await (const tx of client.transactions()) {console.log(`slot=${tx.slot} sig=${tx.signature.toBase58()} matched=${JSON.stringify(tx.filters)}`);}
# pip install decoded-shredstreamfrom decoded_shredstream import Client, Filter, GrpcConfigwith Client.grpc(GrpcConfig(endpoint="<host>:50051", token="<TOKEN>",filters={"all": Filter()})) as client:for update in client:print(update.slot, len(update.data), list(update.filters))
// cargo add decoded-shredstream tokio --features tokio/macros,tokio/rt-multi-threaduse decoded_shredstream::{Filter, GrpcClient, GrpcConfig};let mut client = GrpcClient::connect(GrpcConfig::new("<host>:50051", "<TOKEN>").filter("all", Filter::all()),).await?;while let Some(update) = client.next_update().await {let update = update?;println!("slot={} sig={} matched={:?}", update.slot(), update.signature(), update.filters());}
// go get github.com/shredstream/decoded-shredstream-goclient, err := decodedshredstream.NewGRPC(decodedshredstream.GRPCConfig{Endpoint: "<host>:50051",Token: "<TOKEN>",Filters: decodedshredstream.Filters{"all": decodedshredstream.FilterAll()},})if err != nil { log.Fatal(err) }defer client.Close()err = client.Run(ctx, func(u *decodedshredstream.TransactionUpdate) {sig, _ := u.Signature()fmt.Println(u.Slot, sig, u.Filters)})
Фильтры можно заменять на живом стриме (updateFilters / update_filters) без переподключения. Восстановимые сбои до вас не доходят — клиент переподключается и заново отправляет текущую map фильтров; стрим завершают только отклонённый токен, закрытая сервером сессия или отклонённая map фильтров — через ошибку, которую вы обрабатываете один раз.
💻 Стандартные gRPC-клиенты (сгенерированные стабы)
Для команд, предпочитающих собственный gRPC-стек: контракт — это decoded.proto (наш маршрут, shredstream.com.DecodedShredStreamService), который импортирует свои сообщения из shreder_binary.proto — оба файла поставляются с каждым официальным клиентом. Генерация стабов только из shreder_binary.proto даёт совместимый «из коробки» маршрут, используемый в примерах ниже.
Rust (tonic)
use std::collections::HashMap;use shreder_binary::shreder_binary_service_client::ShrederBinaryServiceClient;use shreder_binary::{SubscribeBinaryTransactionsRequest, SubscribeRequestFilterBinaryTransactions};use solana_transaction::versioned::VersionedTransaction;use tonic::Request;let mut client = ShrederBinaryServiceClient::connect("http://<host>:50051").await?;let filters = HashMap::from([("pumpfun".to_string(),SubscribeRequestFilterBinaryTransactions {account_include: vec!["6EF8rrecthR5Dkzon8Nwu78hRvfCKubJ14M5uBEwF6P".into()],account_exclude: vec![],account_required: vec![],},)]);// Статические фильтры: отправляем один запрос, затем half-close (итератор завершается).let outbound = tokio_stream::iter(vec![SubscribeBinaryTransactionsRequest { transactions: filters }]);let mut req = Request::new(outbound);req.metadata_mut().insert("authorization", "Bearer <TOKEN>".parse()?);let mut stream = client.subscribe_binary_transactions(req).await?.into_inner();while let Some(resp) = stream.message().await? {if let Some(tx) = resp.transaction.and_then(|u| u.transaction) {let vtx: VersionedTransaction = bincode::deserialize(&tx.binary_transaction)?;// → маршрутизируйте по resp.filters, обрабатывайте vtx …}}
Наш эквивалентный маршрут —
DecodedShredStreamServiceClient::subscribe_decoded_transactions, с точно теми же типами запроса/ответа. Чтобы обновлять фильтры на лету, заменитеtokio_stream::iterна канал, который держите открытым, иsend-ите в него новый запрос — переподключение не требуется.
TypeScript (@grpc/grpc-js)
import { Metadata } from "@grpc/grpc-js";// … клиент, сгенерированный из shreder_binary.proto …const meta = new Metadata();meta.set("x-token", "<TOKEN>");const call = client.subscribeBinaryTransactions(meta);call.write({transactions: { pumpfun: { accountInclude: ["6EF8rrec…"], accountExclude: [], accountRequired: [] } },});// call.end(); // half-close, если фильтры не меняютсяcall.on("data", (resp) => {const raw = resp.transaction?.transaction?.binaryTransaction; // Buffer// десериализуйте VersionedTransaction на стороне клиента …});call.on("error", (e) => { /* DATA_LOSS / UNAUTHENTICATED → переподключение / лог */ });
Python (grpcio)
import grpcimport shreder_binary_pb2 as pb, shreder_binary_pb2_grpc as rpcchan = grpc.insecure_channel("<host>:50051")stub = rpc.ShrederBinaryServiceStub(chan)def requests():yield pb.SubscribeBinaryTransactionsRequest(transactions={"pumpfun": pb.SubscribeRequestFilterBinaryTransactions(account_include=["6EF8rrec…"])})md = (("x-token", "<TOKEN>"),)for resp in stub.SubscribeBinaryTransactions(requests(), metadata=md):raw = resp.transaction.transaction.binary_transaction # байты bincode# десериализуйте VersionedTransaction …
Быстрый тест с grpcurl
grpcurl -plaintext -proto shreder_binary.proto \-H 'x-token: <TOKEN>' \-d '{"transactions":{"all":{}}}' \<host>:50051 shreder_binary.ShrederBinaryService/SubscribeBinaryTransactions
🧯 Обработка ошибок и переподключение
| Статус gRPC | Значение | Что делать |
|---|---|---|
UNAUTHENTICATED | токен отсутствует / недействителен / не тот поток | исправьте токен или поток — не повторяйте вслепую |
INVALID_ARGUMENT | некорректный base58, слишком много pubkeys, слишком много именованных фильтров | исправьте проблемный фильтр |
DATA_LOSS | ваш потребитель слишком медленный — очередь отправки сервера переполнилась и стрим закрыт | переподключитесь и заново отправьте вашу map фильтров |
PERMISSION_DENIED | отключение/отзыв оператором | не зацикливайтесь; свяжитесь с оператором |
Сервис работает только в реальном времени — без повтора и дозагрузки. При DATA_LOSS или обрыве транспорта переподключайтесь с экспоненциальным backoff (и небольшим джиттером), затем заново отправьте вашу map фильтров; короткий пропуск данных во время переподключения — это нормально. Из-за правила «одно соединение на токен» убедитесь, что на токен работает только один цикл переподключения.
Чтобы вообще не допускать DATA_LOSS, вычитывайте стрим не блокируясь: передавайте тяжёлую обработку (включая десериализацию bincode) в очередь или пул воркеров, чтобы ваш цикл приёма никогда не застревал. gRPC поверх HTTP/2 гарантирует порядок и целостность всего отправленного — единственная возможная потеря — это сброс «клиент слишком медленный», и он всегда явно сигнализируется через DATA_LOSS.
⚖️ gRPC или UDP?
| gRPC | UDP | |
|---|---|---|
| Задержка | Низкая — h2c, без рукопожатия TLS | Минимальная — бинарные датаграммы, без соединения |
| Доставка | Упорядоченная, надёжная; «клиент слишком медленный» сигнализируется через DATA_LOSS | Push по принципу best-effort; потери возможны, обнаруживаются по пропускам seq |
| Фильтрация | Фильтры по аккаунтам на стороне сервера (статические ключи) | Нет — фильтруйте на стороне клиента после получения |
| Доступность | Исходящее соединение — работает за NAT/файрволами | Требует публично доступной конечной точки UDP |
| Полезная нагрузка | VersionedTransaction (bincode) | VersionedTransaction (bincode), в конверте из 16 байт |
Выбирайте gRPC, когда нужны серверная фильтрация, упорядоченная доставка и совместимость с NAT. Выбирайте UDP ради абсолютно минимальной задержки, если у вас есть доступная конечная точка и вы сами отслеживаете потери. Полезная нагрузка в обоих случаях идентична.
➡️ Дальнейшие шаги
- Приём декодированных транзакций — UDP-конверт, смещения, обнаружение пропусков и минимальный декодер.
- Decoded Shred Stream — позиционирование, задержка и режимы доставки.